The Wayback Machine - https://web.archive.org/web/20160205222236/http://blog.devork.be/search/label/python

devork

E pur si muove

Pylint and dynamically populated packages

Thursday, December 04, 2014

Python links the module namespace directly to the layout of the source locations on the filesystem. And this is mostly fine, certainly for applications. For libraries sometimes one might want to control the toplevel namespace or API more tightly. This also is mostly fine as one can just use private modules inside a package and import the relevant objects into the __init__.py file, optionally even setting __all__. As I said, this is mostly fine, if sometimes a bit ugly.

However sometimes you have a library which may be loading a particular backend or platforms support at runtime. An example of this is the Python zmq package. The apipkg module is also a very nice way of controlling your toplevel namespace more flexibly. Problem is once you start using one of these things Pylint no longer knows which objects your package provides in it's namespace and will issue warnings about using non-existing things.

Turns out it is not too hard to write a plugin for Pylint which takes care of this. One just has to build the right AST nodes in place where they would be appearing at runtime. Luckily the tools to do this easily are provided:

def transform(mod):
    if mod.name == 'zmq':
        module = importlib.import_module(mod.name)
        for name, obj in vars(module).copy().items():
            if (name in mod.locals or
                    not hasattr(obj, '__module__') or
                    not hasattr(obj, '__name__')):
                continue
            if isinstance(obj, types.ModuleType):
                ast_node = [astroid.MANAGER.ast_from_module(obj)]
            else:
                if hasattr(astroid.MANAGER, 'extension_package_whitelist'):
                    astroid.MANAGER.extension_package_whitelist.add(
                        obj.__module__)
                real_mod = astroid.MANAGER.ast_from_module_name(obj.__module__)
                ast_node = real_mod.getattr(obj.__name__)
                for node in ast_node:
                    fix_linenos(node)
            mod.locals[name] = ast_node

As you can see the hard work of knowing what AST nodes to generate is all done in the astroid.MANAGER.ast_from_module() and astroid.MANAGER.ast_from_module_name() calls. All that is left to do is add these new AST nodes to the module's globals/locals (they are the same thing for a module).

You may also notice the fix_linenos() call. This is a small helper needed when running on Python 3 and importing C modules (like for zmq). The reason is that Pylint tries to sort by line numbers, but for C code they are None and in Python 2 None and an integer can be happily compared but in Python 3 that is no longer the case. So this small helper simply sets all unknown line numbers to 0:

def fix_linenos(node):
    if node.fromlineno is None:
        node.fromlineno = 0
    for child in node.get_children():
        fix_linenos(child)

Lastly when writing this into a plugin for Pylint you'll want to register the transformation you just wrote:

def register(linter):
    astroid.MANAGER.register_transform(astroid.Module, transform)

And that's all that's needed to make Pylint work fine with dynamically populated package namespaces. I've tried this on zmq as well as on a package using apipkg and its seems to work fine on both Python 2 and Python 3. Writing Pylint plugins seems not too hard!

New pytest-timeout release

Thursday, August 07, 2014

At long last I have updated my pytest-timeout plugin. pytest-timeout is a plugin to py.test which will interrupt tests which are taking longer then a set time and dump the stack traces of all threads. This was initially developed in order to debug some some tests which would occasionally hang on a CI server and can be used in a variety of similar situations where getting some output is more useful then getting a clean testrun.

The main new feature of this release is that the plugin now finally works nicely with the --pdb option from py.test. When using this option the timeout plugin will now no longer interrupt the interactive pdb session after the given timeout.

Secondly this release fixes an important bug which meant that a timeout in the finaliser of a fixture at the end of the session would not be caught by the plugin. This was mainly because pytest-timeout was not updated since py.test changed the way fixtures where cached on their scope, the introduction of @pytest.fixture(scope='...'), even though this was a long time ago.

So if you use py.test and a CI server I suggest now is as good a time as any to configure it to use pytest-timeout, using a fairly large timeout of say 300 seconds, then forget about it forever. Until maybe one day it will suddenly save you a lot of head scratching and time.

Designing binary/text APIs in a polygot py2/py3 world

Sunday, April 27, 2014

The general advice for handling text in an application is to use a so called unicode sandwich: that is decode bytes to unicode (text) as soon as receiving it, have everything internally handle unicode and then right at the boundary encode it back to bytes. Typically the boundaries where the decoding and encoding happens is when reading from or writing to files, when sending data across the network etc. So far so good.

All this is fine in an environment where it is possible to know the encoding to be used and where an encoding failure can simply be treated as a hard failure. However POSIX is notoriously bad at this, for many things the kernel just doesn't care and any bytes which go in will come back out. This means that for e.g. a filename or command line arguments the kernel does not care about it being valid in the current locale/encoding or even any encoding. When Python 3.0 was initially released this was a problem and by Python 3.1 the solution used was to introduce the surrogateescape error handler for decoders and encoders. This allows Python 3 to smuggle un-decodable bytes in unicode strings and the encoder will put them back when round-tripping. The classical example of why this is useful is when listing files using e.g. os.listdir() to then later pass them back to the kernel via e.g. open().

The downside of surrogate escapes is that the unicode strings now are no longer valid for many other normal string manipulations. If you try to write the result of os.listdir() to a file which you want to encode using UTF8 the encoding step will blow up, so this kind of brings the old Python 2 situation with bytes back. So any user of the API needs to be aware that strings may contain surrogate escapes and handle them appropriately. For a detailed description of these cases refer to Armin Ronacher's Unicode guide which introduces is_surrogate_escaped(s) and remove_surrogate_escaping(s, method='ignore') functions which are pretty self-explanatory.

But let's for now accept the surrogate escape solution Python 3 introduces, as long as the API documents this a user can handle it with the earlier mentioned helper functions. However when designing a polygot library API it is impossible to use the surrogateescape error handler since it does not exist for Python 2.7. And since the required groundwork was not backported either it is impossible to write a surrogateescape handler for Python 2.7, which I consider a glaring omission certainly given the timeline. So this pretty much makes the surrogateescape option not viable as a 2.7/3.x API.

So what options are there left for an API designer? One suggestion is to use native strings: bytes on Python 2.7 and unicode with surrogateescapes on Python 3.x. This means in either case there is no loss of data. But it also means the user of the API now has a harder time writing polygot code if they want to use the unicode sandwich. Given the difficulties to the user I'm not sure I'm a fan of this API.

Another correct, but rather unfriendly, option is to just consider the API to expose bytes and provide the encoding which should be used to decode it. In this case the user can choose the appropriate error handler themselves, be it =ignore=, =replace= or, on Python 3, surrogateescape. The advantage is that this would behave exactly the same on Python 2 and Python 3, however it leaves a casual user of the API a bit lost, certainly on Python 3 where receiving bytes of the API is not very friendly and feels like pushing the Python 2 problems back onto them.

Yet another option I've been considering is provide both APIs: one exposing the bytes, with the attributes possibly prefixed with a b, and one convenience API which decoded the bytes to unicode using the =ignore= error handler. This really seems to pollute the API but might still be the most pragmatic solution: it behaves the same on both Python 2 and Python 3, does not lose any information, allows easy use of the all-unicode inside text model yet still allows explicit handling of decoding.

So what is the best way to design a polygot API? I would really like to hear peoples opinions on which API would be the nicest to use. Or hear if there are any other tricks to employ for polygot APIs.

Saturday, February 08, 2014

Don't be scared of copyright

It appears there is some arguments against putting copyright statements on the top of a file in free software or open source projects. Over at opensource.com Rich Bowen argues that it is counter productive and not in the community spirit (in this case talking about OpenStack). It seems to me the main arguments are the following:

  • It is intimidating to new people
  • It gets too verbose
  • Encourages contribution for the wrong reasons
  • It is hard to decide when to add a name
  • It is even harder to decide when to remove a name
  • The VCS keeps track of contributions anyway

Lastly and perhaps most improtantly he asks:

[...] why do you care? What are you trying to protect against? If you're trying to protect against your contribution being taken by the community and used for other purposes, perhaps contributing to an Apache-licensed code base isn't the smartest thing to do.

Now I think the last question is the most important to answer: you want to assert your copyright on a file to avoid your work from being re-licenced against your will.

That to me is really the crux of the issue and it certainly does not go against the spirit of free or open source software. In fact every additional author asserting their copyright under the license chosen for the project makes the committment of the project to this free or open source license even deeper. It is an attempt to protect against some hypotetical future lawyer who might one day try to claim someone was allowed to do something which was not in the spirit of the free or open source project. For every additional person or organisation listed as holding copyright it becomes harder to ever re-license the work. And this is a good thing.

Now I choose the words "assert your copyright" carefully. Not being a lawyer I am not sure what is the best way to do this. Personally I remember the FSF recommending to put a copyright line with the license in every file so I trust them that this is a fairly legally sound approach. But likewise I'm fine with a LICENSE.txt and AUTHORS.txt file, it is far nicer to work with as developers however I can imagine lawyers being able to attack such a system easier.

As to addressing the other minor points: on the social issues I can't really counter much. Yes it would be a shame if people would be scared away for no reason.

As for deciding when to add or remove a person to the copyright: this might always remain tricky, but if in doubt add the person. And simply never ever remove a person. However the claim that the VCS would be capable of tracking who is the author of what fragment is a bit frivolous, it is fairly common for patches to be applied by a committer instead of the author. And even if originating from a pull request there might easily be reasons to change the commit for some typos, squashing commits etc. I've seen the original authors of a commit disappear before, so really wouldn't just leave it up the VCS to do this bookkeeping.

So in short, the more people are listed as owning copyright on a project the healthier it is and the more I trust it. Please do not be scared away by other people or organisations being listed as copyright holders.

Setting up a Python development environment on Solaris 11

Sunday, December 08, 2013

Solaris is one of those things which tends to be a bit un-loved, it being baught by Oracle probably didn't help either. But regardless they've made a rather nice and modern OS out of it and it is used in a fair few places so you may need to support it. Here I'll describe how to easily setup a Python development environment on Solaris 11 x86. This will use Vagrant and VirtualBox for virtualisation and bootstrap the development environment using the excellent OpenCSW packages.

Building a Solaris Vagrant box

Vagrant is a fairly nice tool to manage local development virtual machines. It's automation of the ssh setup and project folder sharing are very nice. However all docs encourage you to find a "box" image somewhere on the net and simply update it. This won't work for Solaris as there's no box image, and I'm not providing one either as it'd probably break the Solaris license. Also, if you created the box from scratch at least you know what's in it. There are a few tools that try to automate building boxes, but they're more layers of complexity for little gain (at least for this scenario). And as ususal you can bet their Solaris support is much less smooth then say their Ubuntu support.

There's very little information on how to build a box from scratch however. But it turns out it's not that hard. Before you get started however make sure you have Vagrant 1.3 or above as Vagrant isn't quite completely OS-agnostic and learned how to handle Solaris in 1.3.

Firstly, download the install ISO from Oracle, the text installer is my preferred one as there is really no need for all the UI stuff which just takes space in your final box. I'm not providing a direct link as Oracle URLs change more often then the phase of the moon, but you should be able to find the Solaris 11 download links somewhere. There should even be no need to sign up or register at all, just agree to the license IIRC.

Now create a new VirtualBox VM as normal, but when walking through the installer setup screens there's a few points to observe:

  • Use something like "vagrant-solaris11" as hostname. I'm not sure how important this is but I think it's a vagrant convention.
  • Solaris has default password restrictions. Something like "vagrant1" worked for me but just make something up, we'll change it later.
  • Use "vagrant" as the username of the primary account. Each box needs this user.

Once installed a few things need to be setup for vagrant to fully work. This is all insecure, but this is a local box anyway so no need to worry (right?).

  • Edit /etc/default/passwd to set NAMECHECK=NO and MINNONALPHA=0. This will allow us to create insecure passwords.
  • Now change the root password to "vagrant" using passwd, do the same for the "vagrant" account.
  • Allow passwordless use of sudo: visudo -f /etc/sudoers.d/svc-system-config-user and edit this so that the entry looks like: vagrant ALL=(ALL)NOPASSWD: ALL.
  • Vagrant uses a known ssh key to log into the VMs. Again, this is insecure so do not expose this box to anything but yourself. Do something like this as the vagrant user to install this key:
            $ cd ~
            $ mkdir .ssh
            $ wget -O .ssh/authorized_keys \
               https://raw.github.com/mitchellh/vagrant/master/keys/vagrant.pub
            $ chmod 700 .ssh
            $ chmod 600 .ssh/authoriszed_keys
        

Finally you need to install the VirtualBox guest additions. After clicking the menu item to install them Solaris should auto-mount the guest additions ISO somewhere in /media/VBOXADDITIONS_.... Install the additions with

    pkgadd -G -d /media/VBOXADDITIONS_.../VBoxSolarisAddtions.pkg

Time to create a vagrant box out of this virtual machine. Shut it down and find either it's name or UUID using vboxmanage list vms. Now you can use vagrant to build the box:

    vagrant package --base $name_or_uuid --output solaris11.box

This was a fair amount of manual work, but the good news is that you'll never have to do this again. Save this box somewhere safe and it can be used as the base of any new boxes you want for Solaris 11.

Creating the Python development environment

Now the manual way of installing a Python environment is somehow install a compiler, hunt down all Python's dependencies, build them from source and then build Python. Followed by figuring out what is still missing and repeating the whole process.

Fortunately the good people at OpenCSW have created an entire APT-like package repository based on the GNU toolchain. And even better then that they are pretty good at keeping it updated. So in order to setup a Python environment we just need to install a few of their packages. And even if you wanted to compile Python from source I'd reccomend to use OpenCSW packages for the dependencies.

There are detailed instructions on OpenCSW's site, but I'll continue our earlier example here.

  • First create a new Vagrant project:

            $ mkdir myproj
            $ cd myproj
            $ vagrant init solaris11 /path/to/solaris11.box
        

    This creates a Vagrantfile in the current directory, which is the per-project Vagrant configuration. You can have a look, but mostly don't need to change anything for this simple scenario.

  • Time to startup the box and ssh into it:

            $ vagrant up
            ...
            $ vagrant ssh
        
  • Once the virtual machine is running and you are logged into it it's time to install some more packages. Firstly bootstrap OpenCSW:

            # pkgadd -d http://get.opencsw.org/now
        

    (Note how this is utterly insecure)

  • This installed pkgutil which is a package manger for the old SVR4 packages, that is they are the dpkg/rpm equivalent and pkgutil was an apt/yum layer on top of that which Solaris was missing before Solaris 11. OpenCSW is working on supporting Solaris 11 with the native IPS system, but until then pkgutil still works a treat.

    Before doing anything else you want to configure the mirror and release stream you want to use. Edit /etc/opt/csw/pkgutil.conf to set the mirror to http://mirror.opencsw.org/opencsw/unstable or any other unstable mirror. Other release streams will have a lot less recent software available and this is a development box after all

  • Now I prefer to setup at least some form of package verification. So lets set up the key infrastucture for this:

            # pkgutil -U
            # pkgutil -i cswpki
        

    Once this is finished you'll want to enable this in the configuration. Edit /etc/opt/csw/pkgutil.conf again but this time set use_gpg=true and use_md5=true. Unless you like seeing warnings you may also want to configure the OpenCSW key to be ultimately trusted. It's not like you use this GPG account for anything else.

Now you just need to install the required packages using pkgutil. The following sets you up with Python 2.7:

    # pkgutil -i python27 py_setuptools py_pip

If you want you can also install python33 for Python 3.

If you'd also like to play with extension modules you need a few more things. This is the package list needed for using cffi:

  • gcc4core
  • python27_dev
  • libffi_dev

However before you can install cffi there's one more thing you need to do. The default installation when using the text installer is missing some system headers you need. Luckly since Solaris 11 there is the very powerful Image Packagin System (IPS) to manage Solaris packages (OpenCSW is still working on supporting this natively, pkgutil was made at a time IPS did not exist in Solaris). This is a rather impressive package manager if you feel like digging into it, it's kind of to apt/yum like what ZFS was to ext3 or like SMF was to sysv-init, if that makes sense. And what's even better is that Oracle is hosting the basic Solaris 11 package repository for free so you can use it without any extra setup to install the missing headers just like you've been able to do with apt for the last 20 years:

    pkg install [-nv] system/header

Here the -nv options are to show you a preview before actually doing things, which is a good habit to get into I guess.

And at this point you should be able to install cffi using pip just like you're used too: pip install --user cffi.

Re-packaging the vagrant box

Once you have setup the development environment you want it might be worth to re-package the box again before you start on your project. This will allow you to use it as the base of another project later. It will also allow you to simply throw away the box using vagrant destroy if you managed to somehow break it after which a simple vagrant up would initialise you a new VM with the python environment ready to use.

This is very similar as creating the first box. The only difference is that you want to run vagrant in the directory of your project and not specify the base. This will make vagrant use the VM in the project as the base of the new box:

    vagrant package --output /path/to/boxes/dir/solaris11-py27.box

Now the last thing you may want to do is update the Vagrantfile of you project to point to this new box.

Inspecting un-imported modules using ast

Sunday, June 24, 2012

or

Skipping modules in py.test using marks but no importing

As I've said before, py.test is my favourite testing tool. One if it's features is that it allows you to mark tests with arbitrary "marks", built in ones are e.g. skipif, xfail etc. But py.test's extension mechanism allows you to easily add behaviour on any other marks you might want.

Recently I've been starting to write some testing code for a Django project (using the great-and-still-improving pytest-django plugin), but I have one strange requirement, on our CI hosts we only want to run the Django tests for a few platforms (We write monitoring software so support way more platforms then is sane). So the natural thing to was mark those test modules as needing Django so we can skip them.

This is how you mark a test module as such, each test written in this module will have the "django" mark applied:

# test_foo.py
import django
import pytest
import foo

pytestmark = pytest.mark.django
# Or, using multiple marks
pytestmark = [pytest.mark.django, py.test.mark.another_mark]

I've also imported django itself, usually test modules need to do this one way or another, at least indirectly via the module they are testing. Now we could easily write py.test conftest file to skip these tests:

# conftest.py
import pytest

ENABLE_DJANGO = True  # imagine some complicated expression/function

def pytest_setup_item(item):
    if 'django' in item.keywords and not ENABLE_DJANGO:
        pytest.skip('Django tests are disabled')

But we have an additional problem, we don't even build Django for all platforms so can't import these test modules in the first place. This means py.test fails while collecting the test modules. There is another hook available in py.test which allows us to skip modules without importing them: pytest_ignore_collect(). But we still need to be able to read the marker inside this module without importing, time to enter ast.

The ast module can parse python code into an abstract syntax tree, which represents the code in a tree of nodes according to the grammar. This is actually a part of the normal compilation process, but just as an internal step. Enough introduction, it was pleasantly easy to use it to find the marks in our test modules:

def _static_mark(path):
    """Return the pytestmark of a test file whithout executing it"""
    module = ast.parse(path.read(), str(path))
    marks = []
    for node in module.body:
        if not isinstance(node, ast.Assign):
            continue
        if len(node.targets) > 1:
            continue
        if not isinstance(node.targets[0], ast.Name):
            continue
        if node.targets[0].id != 'pytestmark':
            continue
        if isinstance(node.value, ast.Attribute):
            marks.append(node.value.attr)
            continue
        if isinstance(node.value, ast.List):
            for elem in node.value.elts:
                if isinstance(elem, ast.Attribute):
                    marks.append(elem.attr)
    return marks

That's it. Essentially we're looking for an ast.Assign node which has a target of pytestmark that is, something is being assigned to pytestmark. Once we find such a node we make sure to only accept a few right-hand-side expressions, namely either pytest.mark.the_mark or [pytest.mark.mark0, pytest.mark.mark1].

Now this is obviously not bulletproof, but it does keep it nice and straight forward. And I thought is a nice example of how to the ast module can be useful.

Small Emacs tweaks impoving my Python coding

Thursday, September 15, 2011

Today I've spent a few minutes tweaking Emacs a little. The result is very simple yet makes a decent impact on usage.

Firstly I remembered using the c-subword-mode a long time ago, I couldn't believe I never used that in Python. Turns out there is a more genericly named subword-mode by now (emacs 23 IIUC) and it's very easy to enable for Python by default:

(add-hook 'python-mode-hook (lambda () (subword-mode 1)))

The second simple improvement was finally figuring out how to automatically enable flyspell-mode when editing Restructured Text:

(add-hook 'rst-mode-hook (lambda () (flyspell-mode 1)))

Both where so simple I can't believe I didn't do them earlier.

And while on the subject, develock-mode is great and I've been using it a very long time. Unfortunately python isn't supported out of the box, but no worries, someone has done the work already. So all you have to do is something along the lines of:

(load "~/.emacs.d/develock-py.el")

I don't think I'd still want to edit a python file without it

Using __getattr__ and property

Friday, June 17, 2011

Today I wasted a lot of time trying to figure out why a class using both a __getattr__ and a property mysteriously failed. The short version is: Make sure you don't raise an AttributeError in the property.fget()

The start point of this horrible voyage was a class which looked roughly like this:

class Foo(object):

    def __getattr__(self, name):
        return self._container.get(name, 'some_default')

    @property
    def foo(self):
        val = self._container.get(foo)
        if test(val):
            return some_helper(val)
        return val

This sees fine enough. Only it turns out that some_helper() raised an AttributeError for some invalid input. Certainly reasonable since it was never meant to deal with incorrect input, that was meant to have been sanitised already (that was a bug in the caller which was actually just an incorrect unittest). The main gotcha was that it seems that python doesn't just check whether a "foo" is present in all the relevant __dict__'s along the mro. Instead it seems to use getattr(inst, "foo") and then delegate to __getattr__() if it gets an AttributeError. Now suddenly finding a bug in some_helper() has turned into a puzzling question as to why __getattr__() was called.

Personally I can't see why it doesn't use the mro to statically look up the required object instead of using the AttributeError-swallowing approach. But maybe there's a good reason.

Synchronising eventlets and threads

Thursday, March 17, 2011

Eventlet is an asynchronous network I/O framework which combines an event loop with greenlet-based coroutines to provide a familiar blocking-like API to the developer. One of the reasons I like eventlet a lot is that the technology it builds on allows it's event loop to run inside a thread or even run multiple event loops in different threads. This makes it a lot more amenable to slowly evolving existing applications to be more asynchronous then solutions like e.g. gevent which only allow one event loop per process.

Eventlet isn't the most mature of tools however and it's API shows signs of being developed as needs arose. Not that this is necessarily a bad thing, APIs do need to grow from being used, but don't be surprised if you need to dig down into some parts and discover rough edges (hi IPv6!). The API does have a decent collection of tools you'll be familiar with however: greenlet-local storage, semaphores, events (though not quite the event you're used too) and even some extra goodies like pools, WSGI servers, DBAPI2 connection pools and ZeroMQ support. But you'll notice that all these goodies are designed to work in a greenlet-only world. And the one place where threads are acknowledged, a global threadpool of workers as a last resort to make things behave asynchronous, looks very messy and entirely not reusable (it's full of module globals for one). So if you're introducing eventlet into an existing application and you need to communicate data and events between threads and eventlets you'll find a void.

So after some studying of the tpool module I decided to build a class which could synchronise between threads and eventlets. This class, which I called a Notifier, can be basically thought of as a Condition without the lock, i.e. there are three methods: .wait(), .notify() and the rather similar .notify_all(). The idea is that any thread or eventlet which calls .wait() will block (cooperatively block in the case of a greenlet) until it gets woken up by a call to one of the notifying methods. That's all there is to it.

Building a Notifier

(Be prepared to look at the source code for eventlet.hubs.hub, threading and related code when reading this.)

Firstly the class will need to be constructed. For now there's only one interesting instance attribute and that is _waiters which is a set which will contain all the threads and eventlets currently blocking on a call to .wait().

def __init__(self, hubcache=GLOBAL_HUBCACHE):
        self._waiters = set()
        self.hubcache = hubcache
Don't worry yet about what goes into the set of waiters and also ignore the hubcache for now. We'll get to those later.

Now lets build the .wait() method: it needs to block until notified. But blocking is significantly different when you're running in a thread then when running in an eventlet. The basics are that in a thread you want to really block using the locking primitives provided by the OS (exposed to Python in the thread module) while an eventlet essentially wants to switch to the event loop, called the hub, in the hope someone will eventually switch back to it when it needs to wake up. These two are so different that they easily divide in two methods: .gwait() for blocking eventlets and .twait() for blocking threads.

def wait(self, timeout=None):
    hub = eventlet.hubs.get_hub()
    if hub.running:
        self.gwait(timeout)
    else:
        self.twait(timeout)
You could call these directly obviously, but as you can see abstracting them away is not that hard. You can easily detect if you're in an eventlet by checking if the hub (which is thread-local) is actually running. The price you pay for this is that this will create a hub instance in each thread, even if it is not used. But the worst this does is waste some memory. (This could be avoided by checking for the eventlet.hubs._threadlocal.hub attribute, but that's even more internal then .get_hub().)

As already mentioned the basics of blocking in an eventlet is to switch to the hub and then wait until some other greenlet running in the same thread switches back to you. So a notifier needs to have a reference to your geenlet instance so it can call .switch() on it when the time comes to wake you up. But what does another thread do? Well it turns out another thread could ask your hub to schedule a function to run in your thread using the hub's .schedule_call_global() method since the only thread-critical operation is an append on the hubs's next_timers list, which is a thread-safe operation. Now there is another catch, remember that the hub is basically an event loop? Well if no events happen then it will not be going round it loop! And appening something to a list is not creating an event. So what you need to do is use os.pipe so you have a filedescriptor which you can register with the hub and now you can just write some data into this pipe when you want the hub of the waiter to wake up. Setting up this pipe is what the mysterious call to ._create_pipe() does, we'll see it in detail later. The rest is just simple sugar: dealing with timeouts and returning the correct values for them:

def gwait(self, timeout=None):
    waiter = eventlet.getcurrent()
    hub = eventlet.hubs.get_hub()
    self._create_pipe(hub)
    self._waiters.add((waiter, hub))
    if timeout and timeout > 0:
        timeout = eventlet.Timeout(timeout)
        try:
            with timeout:
                hub.switch()
        except eventlet.Timeout, t:
            if t is not timeout:
                raise
            self._waiters.discard((waiter, hub))
            return False
        else:
            return True
    else:
        hub.switch()
        return True

Next on lets look at how you block in a thread. This is actually surprisingly simple, just copy what threading.Condition.wait() does for it: it is perfectly non-blocking for a greenlet to call .release() on a lock which was acquired by another thread. Notice the ugly CPU-consuming spinning which happens when a timeout is in use, luckily this has been fixed in python 3.2 (issue7316).

def twait(self, timeout=None):
    waiter = threading.Lock()
    waiter.acquire()
    self._waiters.add((waiter, None))
    if timeout is None:
        waiter.acquire()
        return True
    else:
        # Spin around a little, just like the stdlib does
        _time = time.time
        _sleep = time.sleep
        min = __builtin__.min
        endtime = _time() + timeout
        delay = 0.0005      # 500 us -> initial delay of 1 ms
        while True:
            gotit = waiter.acquire(0)
            if gotit:
                break
            remaining = endtime - _time()
            if remaining <= 0:
                break
            delay = min(delay * 2, remaining, .05)
            _sleep(delay)
        if not gotit:
            self._waiters.discard((waiter, None))
            return False
        else:
            return True
The last thing of interest here is how the hub does not matter, so we just place None in the set of waiters.

Now lets have a look at what notify looks like. We've already discussed what needs to happen here: to notify an eventlet we need schedule a call with it's hub to switch to it (no point special casing when we're already running in that same hub). In case we notified it from a different thread we also need to signal the hub using the pipe set up by .wait() so it will actually start to go round it's loop and execute this just scheduled call. Notifying a thread is even easier: just unlock the lock it's trying to acquire.

def notify(self):
    if self._waiters:
        waiter, hub = self._waiters.pop()
        if hub is None:
            # This is a waiting thread
            try:
                waiter.release()
            except thread.error:
                pass
        else:
            # This is a waiting greenlet
            def notif(waiter):
                waiter.switch()
            hub.schedule_call_global(0, notif, waiter)
            if hub is not eventlet.hubs.get_hub():
                self._kick_hub(hub)

Admittedly I've hidden some of the cute trickery away in those ._create_pipe() and ._kick_hub() calls, so lets leave the boring .notify_all() and skip straight to them.

The principle for this is that each eventlet which calls .gwait() needs to ensure there is a pipe available to which a thread can write something. If the hub of the eventlet was waiting for the reading end of this pipe to become readable it will wake up and notice it has to run the notif() function which was scheduled by our call to .notify(). But creating a pipe for each call to .gwait() does seem rather wasteful and this is where the mysterious hubcache comes into play: it is a dictionary keeping track of the pipe associated with each hub.

def _create_pipe(self, hub):
    if hub in self.hubcache:
        return
    def read_callback(fd):
        os.read(fd, 512)
    rfd, wfd = os.pipe()
    listener = hub.add(eventlet.hubs.hub.READ, rfd, read_callback)
    self.hubcache[hub] = (rfd, wfd, listener)
You can see this asks the hub to wake up when the reading end of the created pipe becomes readable and when this happens read_callback() will be called. The only purpose of read_callback() is to read all the data written to the pipe so that the OS buffers are emptied and the pipe can be re-used to wake the hub up the next time.

Now there is just ._kick_hub() left. This should now be obvious: look up the hub in the hubcache and write some data to the writing end of the pipe. The only gotcha here is that while we might be called from another thread, this does not mean the calling thread itself can't be part of an eventlet mainloop. So in that case make sure not to do a blocking write (as unlikely as that might be).

def _kick_hub(self, hub):
    rfd, wfd, r_listener = self.hubcache[hub]
    current_hub = eventlet.hubs.get_hub()
    if current_hub.running:
        def write(fd):
            os.write(fd, 'A')
            current_hub.remove(w_listener)
        w_listener = current_hub.add(eventlet.hubs.hub.WRITE, wfd, write)
    else:
        os.write(wfd, 'A')
Notice here how in the async case we remove the listener which was used to wake the hub up right from the callback itself. If we didn't do this then the hub would most likely find the writing end of the pipe writable again on the next loop which would trigger another notification of the other hub, not what we want!

That's cute, but what now?

We've now got a great way of notifying other threads and eventlets at will. But this is an entirely non-standard tool! Using this is strange, unfamiliar and unwieldy. This isn't one of the synchronisation primitives we know and wanted to use.

But look how easy it is to build a lock now: we just need a real lock and one of these strange notifiers:

class Lock(object):

    def __init__(self, hubcache=GLOBAL_HUBCACHE):
        self.hubcache = hubcache
        self._notif = Notifier(hubcache)
        self._lock = threading.Lock()
        self.owner = None

    def acquire(self, blocking=True, timeout=None):
        gotit = self._lock.acquire(False)
        if gotit or not blocking:
            if gotit:
                self.owner = eventlet.getcurrent()
            return gotit
        if timeout is None:
            while not gotit:
                self._notif.wait()
                gotit = self._lock.acquire(False)
            if gotit:
                self.owner = eventlet.getcurrent()
                return True
        else:
            if timeout < 0:
                raise RuntimeError('timeout must be greater or equal then 0')
            now = time.time()
            end = now + timeout
            while not gotit and (now < end):
                self._notif.wait(end - now)
                gotit = self._lock.acquire(False)
                now = time.time()
            if gotit:
                self.owner = eventlet.getcurrent()
            return gotit

    __enter__ = acquire

    def release(self):
        self._lock.release()
        self.owner = None
        self._notif.notify()

    def __exit__(self, exc_type, exc_value, traceback):
        self.release()

    def __repr__(self):
        return ('<gsync.Lock object at 0x%x (%r, %r)>' %
                (id(self), self.owner, self._notif))
Most of this code is to deal with the timeout of the lock. If you would only implement a lock as how it was before Python 3.2 this would have been even simpler. (Oh, also notice the owner attribute, that will come in handy for the next tool.)

Of course now we have a lock we can easily build the next tool: a proper Condition. It's almost like I planned this!

class Condition(object):

    def __init__(self, lock=None, hubcache=GLOBAL_HUBCACHE):
        if lock is None:
            self._lock = Lock(hubcache=hubcache)
        else:
            self._lock = lock
        self._notif = Notifier(hubcache=hubcache)

        # Export the lock methods
        self.acquire = self._lock.acquire
        self.release = self._lock.release
        self.__enter__ = self._lock.__enter__
        self.__exit__ = self._lock.__exit__

    def __repr__(self):
        return '<gsync.Condition (%r, %r)>' % (self._lock, self._notif)

    def wait(self, timeout=None):
        if self._lock.owner is not eventlet.getcurrent():
            raise RuntimeError('Can not wait on un-acquired lock')
        self._lock.release()
        try:
            self._notif.wait(timeout)
        finally:
            self._lock.acquire()

    def notify(self):
        if self._lock.owner is not eventlet.getcurrent():
            raise RuntimeError('Can not notify on un-acquired lock')
        self._notif.notify()

    def notify_all(self):
        if self._lock.owner is not eventlet.getcurrent():
            raise RuntimeError('Can not notify on un-acquired lock')
        self._notif.notify_all()
No surprises here. This is truly getting trivial to implement thanks to our previous two primitives.

Now once we have a condition we can finally get to the real prise: a queue to move data freely between threads and eventlets.

import Queue as queue

class BaseQueue(object):

    def __init__(self, maxsize=0, hubcache=GLOBAL_HUBCACHE):
        self.hubcache = hubcache
        self.maxsize = maxsize
        self._init(maxsize)
        self.mutex = Lock(hubcache=hubcache)
        self.not_empty = Condition(self.mutex, hubcache=hubcache)
        self.not_full = Condition(self.mutex, hubcache=hubcache)
        self.all_tasks_done = Condition(self.mutex, hubcache=hubcache)
        self.unfinished_tasks = 0

class Queue(BaseQueue, queue.Queue):
    pass
Great, only had to provide a new .__init__() which uses our own locking primitives, everything else of the stdlib Queue class can be re-used.

But why the strange diversion to create a separate BaseQueue class? Well, it makes making the priority and lifo queues very easy:

class PriorityQueue(BaseQueue, queue.PriorityQueue):
    pass

class LifoQueue(BaseQueue, queue.LifoQueue): pass
That's right, mro FTW!

Caveats when mixing threads and eventlets

There is one issue to watch out for: imagine a thread which consumes items from a queue, spawning eventlets to do the work.

class Worker(threading.Thread):
    def __init__(self, inputq):
        self.inputq = inputq

    def run(self):
        eventlet.sleep(0)  # start the hub
        while True:
            item = self.inputq.get()
            if item is PoisonPill:
                break
            eventlet.spawn(self.do_stuff, item)
Notice that eventlet.sleep(0) line? It basically switches to the hub, thereby implicitly starting it, which then switches back immediately. But why?

Remember the code for Notifier.wait(), it checked if the hub was running to detect whether it was being called from inside an eventlet or not. So if you didn't manage to start the hub before calling this method the whole thread will block! Hence the minor hack to start the hub manually beforehand.

All the code

Here all the code for the notifier in one piece, including the docstrings and comments. This also includes the global hubcache with a tiny bit of extra magic to be able to clear the cache if you so desire. Having this as a parameter to pass in allows you to use different hub caches if you have a reason to do so.

import os
import thread
import threading
import time

import eventlet


class HubCache(dict):
    """Cache used by Notifier instances

    This is a dict-subclass to overwrite the .clear() method.  It's
    keys are hubs and values are (rfd, wfd, listener).  Using this
    means you can clear the cache in a way which will unregister the
    listeners from the hubs and close all filedescriptors.

    XXX This is hugely incomplete, only remove items from this cache
        using the .clear() method as the other ways of removing items
        will not release resources properly.
    """

    def clear(self):
        while self:
            hub, (rfd, wfd, listener) = self.popitem()
            hub.remove(listener)
            os.close(rfd)
            os.close(wfd)

    def __del__(self):
        self.clear()


"""The global hubcache

This is the default hubcache used by Notifier instances.
"""
GLOBAL_HUBCACHE = HubCache()


class Notifier(object):
    """Notify one or more waiters

    This is essentially a condition without the lock.  It can be used
    to signal between threads and greenlets at will.
    """
    # This doesn't use eventlet.hubs.trampoline since that results in
    # a filedescriptor per waiting greenlet.  Instead each eventlet
    # that calls .gwait() will ensure there's a filedescriptor
    # registered for reading for with it's hub.  This filedescriptor
    # is then only used when another thread wants to wake up the hub
    # in order for a notification to be delivered to the eventlet.

    def __init__(self, hubcache=GLOBAL_HUBCACHE):
        """Initialise the notifier

        The hubcache is a dictionary which will keep pipes used by the
        notifier so that only ever one pipe gets created per hub.  The
        default is to share this hubcache globally so all notifiers
        use the same pipes for intra-hub communication.
        """
        # Each item in this set is a tuple of (waiter, hub).  For an
        # eventlet the waiter is the greenlet while for a thread it is
        # a lock.  For a thread the hub item is always None.
        self._waiters = set()
        self.hubcache = hubcache

    def wait(self, timeout=None):
        """Wait from a thread or eventlet

        This blocks the current thread/eventlet until it gets woken up
        by a call to .notify() or .notify_all().

        This will automatically dispatch to .gwait() or .twait() as
        needed so that the blocking will be cooperative for greenlets.

        Returns True if this thread/eventlet was notified and False
        when a timeout occurred.
        """
        hub = eventlet.hubs.get_hub()
        if hub.running:
            self.gwait(timeout)
        else:
            self.twait(timeout)

    def gwait(self, timeout=None):
        """Wait from an eventlet

        This cooperatively blocks the current eventlet by switching to
        the hub.  The hub will switch back to this eventlet when it
        gets notified.

        Usually you can just call .wait() which will dispatch to this
        method if you are in an eventlet.

        Returns True if this thread/eventlet was notified and False
        when a timeout occurred.
        """
        waiter = eventlet.getcurrent()
        hub = eventlet.hubs.get_hub()
        self._create_pipe(hub)
        self._waiters.add((waiter, hub))
        if timeout and timeout > 0:
            timeout = eventlet.Timeout(timeout)
            try:
                with timeout:
                    hub.switch()
            except eventlet.Timeout, t:
                if t is not timeout:
                    raise
                self._waiters.discard((waiter, hub))
                return False
            else:
                return True
        else:
            hub.switch()
            return True

    def twait(self, timeout=None):
        """Wait from an thread

        This blocks the current thread by using a conventional lock.

        Usually you can just call .wait() which will dispatch to this
        method if you are in an eventlet.

        Returns True if this thread/eventlet was notified and False
        when a timeout occurred.
        """
        waiter = threading.Lock()
        waiter.acquire()
        self._waiters.add((waiter, None))
        if timeout is None:
            waiter.acquire()
            return True
        else:
            # Spin around a little, just like the stdlib does
            _time = time.time
            _sleep = time.sleep
            min = __builtin__.min
            endtime = _time() + timeout
            delay = 0.0005      # 500 us -> initial delay of 1 ms
            while True:
                gotit = waiter.acquire(0)
                if gotit:
                    break
                remaining = endtime - _time()
                if remaining <= 0:
                    break
                delay = min(delay * 2, remaining, .05)
                _sleep(delay)
            if not gotit:
                self._waiters.discard((waiter, None))
                return False
            else:
                return True

    def notify(self):
        """Notify one waiter

        This will notify one waiter, regardless of whether it is a
        thread or eventlet, resulting in the waiter returning from
        it's .wait() call.

        This will never block itself so can be called from either a
        thread or eventlet itself and will wake up the hub of another
        thread if an eventlet from it is notified.
        """
        if self._waiters:
            waiter, hub = self._waiters.pop()
            if hub is None:
                # This is a waiting thread
                try:
                    waiter.release()
                except thread.error:
                    pass
            else:
                # This is a waiting greenlet
                def notif(waiter):
                    waiter.switch()
                hub.schedule_call_global(0, notif, waiter)
                if hub is not eventlet.hubs.get_hub():
                    self._kick_hub(hub)

    def notify_all(self):
        """Notify all waiters

        Similar to .notify() but will notify all waiters instead of
        just one.
        """
        for i in xrange(len(self._waiters)):
            self.notify()

    def _create_pipe(self, hub):
        """Create a pipe for a hub

        This creates a pipe (read and write fd) and registers it with
        the hub so that ._kick_hub() can use this to signal the hub.

        This keeps a cache of hubs on ``self.hubcache`` so that only
        one pipe is created per hub.  Furthermore this dict is never
        cleared implicitly to avoid creating new sockets all the time.

        This method is always called from .gwait() and therefore can
        only run once for a given hub at the same time.  Thus it is
        threadsave.
        """
        if hub in self.hubcache:
            return
        def read_callback(fd):
            # This just reads the (bogus) data just written to empty
            # the os queues.  The only purpose was to kick the hub
            # round it's loop which is now has.  The notif function
            # scheduled by .notify() will now do it's work.
            os.read(fd, 512)
        rfd, wfd = os.pipe()
        listener = hub.add(eventlet.hubs.hub.READ, rfd, read_callback)
        self.hubcache[hub] = (rfd, wfd, listener)

    def _kick_hub(self, hub):
        """Kick the hub around it's loop

        Threads need to be able to kick a hub around their loop by
        interrupting the sleep.  This is done with the help of a
        filedescriptor to which the thread writes a byte (using this
        method) which will then wake up the hub.
        """
        rfd, wfd, r_listener = self.hubcache[hub]
        current_hub = eventlet.hubs.get_hub()
        if current_hub.running:
            def write(fd):
                os.write(fd, 'A')
                current_hub.remove(w_listener)
            w_listener = current_hub.add(eventlet.hubs.hub.WRITE, wfd, write)
        else:
            os.write(wfd, 'A')

    def __repr__(self):
        return ('<gsync.Notifier object at 0x%x (%d waiters)>' %
                (id(self), len(self._waiters)))

That's all folks

So it seems that with some careful thinkering you can create all your tried and tested tools to communicate between threads and eventlets or between eventlets in different threads. This makes adopting eventlet into an existing application a whole lot more approachable. It certainly helped me!

Creating subprocesses in new contracts on Solaris 10

Wednesday, February 09, 2011

Solaris 10 introduced "contracts" for processes. You can read all about it in the contract(4) manpage but simply put it's a grouping of processes under another ID, and you can "monitor" these groups, e.g. be notified when a process in a group coredumps etc. This is actually one of the tools of Solaris' init replacement smf(5), which is probably the main reason people care about contracts.

Suppose you have a daemon managed by SMF which executes subprocesses as part of it's life (e.g. because of multiprocessing). If your application is not aware of contracts all subprocesses will be in the same contract. Which is fine until the main process dies a miserable death and one expects SMF to restart it. Only it doesn't since it sees the contract as still being alive and sees no harm. Time to start your subprocesses in a different contract then.

Anyway, whatever your motivation, starting the subprocesses in new contracts is normally done by activating the process contract template, i.e. opening /system/contract/process/template and then calling ct_tmpl_activate(3) on the filedescriptor. Only problem is that this isn't exposed to python.

Fortunately the libcontract(3) functions we need for this are very simple, so ctypes can handle this very nicely. The code is pretty simple:

import ctypes
import os

fd = os.open('/system/contract/process/template', os.O_RDWR)
libct = ctypes.cdll.LoadLibrary('libcontract.so.1')
rv = libct.ct_tmpl_activate(fd)
if rv != 0:
    raise Exception('oops')

That's all there is to it. Each fork(2) call will now result in a process running in it's own contract.

It's nice when C-APIs are simple enough to be used from inside python instead of having to write wrappers.

PS: Needless to say you can do a whole lot more with contracts, just read up on the libcontract API docs.

Return inside with statement (updated)

Saturday, August 14, 2010

Somehow my brain seems to think there's a reason not to return inside a with statement, so rather then doing this:

def foo():
    with ctx_manager:
        return bar()

I always do:

def foo():
    with ctx_manager:
         result = bar()
    return result

No idea why nor where I think to have heard/read this. Searching for this brings up absolutely no rationale. So if you know why this is so, or know that the first version is perfectly fine, please enlighten me!

Update:

Seems it's only relevant if you're reading a file using the with statement. This seems to have come from the python documentation itself:

The last version is not very good either — due to implementation details, the file would not be closed when an exception is raised until the handler finishes, and perhaps not at all in non-C implementations (e.g., Jython).

def get_status(file):
    with open(file) as fp:
        return fp.readline()

Sadly it doesn't say what the implementation details are nor how to do this correctly. I can have several more or less educated guesses at why and which way to do this better. But I'd love to get a more detailed description of what happens in the implementations when doing this, because the worst-case way of interpreting that example is that using open() or file() as a context manager is completely useless. Which I would hate.

Templating engine in python stdlib?

Monday, August 09, 2010

I am a great proponent of the python standard library, I love having lots of tools in there and hate having to resort to thirdparty libs. This is why I was wondering if there could be a real templating engine in the stdlib some day? I'm not a great user of templates, but sometimes string.Template is just not good enough and real templating would make things a lot more readable.

Not being a great templating user I don't know how the world of templating looks. Is there a template engine out there which would be a candidate to move to the stdlib? Or do most people think this is a stupid idea? Or maybe this has been discussed before and I didn't find the discussion?

py.test test generators and cached setup

Sunday, June 13, 2010

Recently I've been enjoying py.test's test function arguments, it takes a little getting used too but soon you find that it's quite likely a better way then the xUnit-style of setup/teardown. One slightly more advanced usage was using cached setup together with test generators however. While not difficult that took me some figuring out, so let me document it here.

Since I haven't been a fan of generative tests before I'll explain why I think I can make use of them now. I was writing a wrapper around pysnmp to handle SNMP-GET requests transparently between the different versions. For this I wrote a number of test functions which do some GET requests and check the results, the basic outline of such a test is:

def test_some_get(wrapper_v1):
    oids = ...
    result = wrapper_v1.get(oids)
    assert ...

Here wrapper_v1 is a funcarg which returns an instance of my wrapper class configured for SNMPv1. The extra catch here is that this funcarg uses a function which tries to find an available SNMP agent, trying if one is running on the local host (for the developer) or if a well-know test host is reachable (for lazy developers on our dev network and for buildbots), skipping the test otherwise. But to avoid the relatively long timeouts involved for each individual test this function needs to be cached. Here's the outline of this funcarg:

def pytest_funcarg__wrapper_v1(request):
    cfg = request.cached_setup(setup=check_snmp_v1_avail, scope='session')
    if not cfg:
        py.test.skip('No SNMPv1 agent available')
    return SnmpWrapper(cfg)

Once having all the tests using this wrapper_v1 funcarg I obviously want exactly the same tests for SNMPv2 since that's the whole point of the wrapper. For this I'd need a wrapper_v2 funcarg which is configured for SNMPv2, but that would mean duplicating all the tests! Enter test generators.

The trick to combine test generators with cached setup is not to use the funcargs argument to metafunc.addcall() but rather use the param argument in combination with a normal funcarg. The normal funcarg can then use request.cached_setup() and use the request.param to decide how to configure the wrapper object returned. This is what that looks like:

def pytest_generate_tests(metafunc):
    if 'snmpwrapper' in metafunc.funcargnames:
        metafunc.addcall(id='SNMPv1', param='v1')
        metafunc.addcall(id='SNMPv2', param='v2')

def pytest_funcarg__snmpwrapper(request):
    cfg = request.cached_setup(setup=lambda: check_snmp_v1_avail(request.param),
                               scope='session', extrakey=request.param)
    if not cfg:
        py.test.skip('No SNMP%s agent available' % request.param)
    return SnmpWrapper(cfg)

Don't forget the extrakey argument to cached_setup. The caching uses the name of the requested object, "snmpwrapper" in this case, and the extrakey value to decide when to re-use the caching. If you forget extrakey both calls will return the same cfg.

And that's all that's needed! Test now simply ask for the snmpwrapper funcarg and will get run twice, once configured for SNMPv1 and once for SNMPv2. Running the tests will now look like this:

flub@signy:...$ py.test -v snmp_test.py
============================= test session starts ==============================
python: platform linux2 -- Python 2.6.5 -- pytest-1.1.1 -- /usr/bin/python
test object 1: /home/flub/.../snmp_test.py

snmp_test.py:57: test_get_one.test_get_one[SNMPv1] PASS
snmp_test.py:57: test_get_one.test_get_one[SNMPv2] PASS
snmp_test.py:63: test_get_two.test_get_two[SNMPv1] PASS
snmp_test.py:63: test_get_two.test_get_two[SNMPv2] PASS
snmp_test.py:71: test_get_two_bad.test_get_two_bad[SNMPv1] PASS
snmp_test.py:71: test_get_two_bad.test_get_two_bad[SNMPv2] PASS
snmp_test.py:79: test_get_many.test_get_many[SNMPv1] PASS
snmp_test.py:79: test_get_many.test_get_many[SNMPv2] PASS

=========================== 8 passed in 1.31 seconds ===========================

This wasn't very complicated, but having an example of using the param argument to metafunc.addcall() would have made figuring this out a little easier. So I hope this helps someone else, or at least me at some time in the future.

Update: Originally I forgot the extrakey argument to cached_setup() and thus the funcarg was returning the same in both cases. Somehow I assumed the caching was done on function identity of the setup function. Oops.

Selectable queue

Saturday, May 29, 2010

Sometimes you'd want to use something like select.select() on Queues. If you do some searching for this it turns out that this question has been answered for multiprocessing Queues where you can simply use the real select on the underlying socket (IIRC), but for a good old fashioned Queue you're stuck.

Now it's easy to argue that this isn't that high a need, when I wanted this a while ago it turned out to be surprisingly simple to re-structure the design a little so that I no longer desired a selectable queue. But it's still something that hung around the back of my mind for a while, so I've kept thinking about it. My conclusion for now (which I haven't bothered implementing) is that simply cloning the normal queues but replacing the not_empty and not_full Conditions by Events gives you selectable queues. It changes the semantics slightly, but that doesn't seem harmful. This obviously isn't enough, so the second change to the queue is that you should be allowed to pass in the events to use. And now that you can share the not_full event between two queues you can simply wait on this event and you have your select.

Next time I want this I might actually implement this idea rather then re-design so that I don't want selectable queues anymore.

weakref and circular references: should I really care?

Wednesday, May 12, 2010

While Python has a garbage collector pretty much whenever circular references are touched upon it is advised to use weak references or otherwise break the cycle. But should we really care? I'd like not to, it seem like something the platfrom (python vm) should just provide for me. Are all those mentions really just for the cases where you can't (or don't want to, e.g. embedded) use the normal garbage collector?

Update: By some strange brain failure I seemed to have written "imports" rather then "references" in the title originally. They are obviously a bad thing.

python-prctl

Saturday, May 08, 2010

There is sometimes a need to set the process name from within python, this would allow you to use something like "pkill myapp" rather then the process name of your application just being yet another "python". Bugs (which I'm now failing to find in Python's tracker) have been filed about this and many wannabe implementations made, all of them seemed to have many problems however. Aiming for UNIX cross-compatibility they all messed up and end up trying to make the abstraction everywhere, and ususally the code was less then beautiful

But I've just discovered the python-prctl module by Dennis Kaarsemaker which does something far more sensible: rather then trying to be cross-platform it just wraps the Linux prctl system call (as well as libcap). The code looks well written and the API, while a little bit overloaded on the get-set names, seems nice. It even includes what is probably the most sensible implementation of clobbering argv that I've ever seen (but don't use that, no normal person should ever clobber argv!).

If someone writes a nice module like this to cater for the MacOSX guys, the only other system I know of has a system call to set the process name, then I may never have to worry about getting someting like this into PSI at some distant point in the future. (And speaking about PSI, yes the windows port is still slowly under way. I"m just busy with lots of other things at the same time.)

Storm and sqlite locking

Saturday, May 08, 2010

The Storm ORM struggles with sqlite3's transaction behaviour as they explain in their source code. Looking at the implementation of .raw_execute() the side effect of their solution to this is that they start an explicit transaction on every statement that gets executed. Including SELECT.

This, in turn, sucks big time. If you look at sqlite's locking behviour you will find that it should be possible to read from the database using concurrent connections (i.e. from multiple processes and/or threads at the same time). However since storm explicitly starts that transaction for a select it means the connection doing the read now holds a SHARED lock until you end the transaction (by doing a commit or rollback). But since it's holding this shared lock, for no good reason, it means no other connection can acquire the EXCLUSIVE lock at the same time.

The upshot of this seems to be that you need to call .commit() even after just reading the database, thus ensuring you let go of the shared lock. Can't say I like that.

Python optimisation has surprising side effects

Tuesday, April 27, 2010

Here's something that surprised me:

a = None
def f():
    b = a
    return b
def g():
    b = a
    a = 'foo'
    return b

While f() is perfectly fine, g() raises an UnboundLocalError. This is because Python optimises access to local variables using the LOAD_FAST/STORE_FAST opcode, you can easily see why this is looking at the code objects of those functions:

>>> f.__code__ .co_names
()
>>> f.__code__ .co_varnames 
('a', 'b')
>>> g.__code__ .co_names
('a',)
>>> g.__code__ .co_varnames 
('b',)

I actually found out this difference thanks to finally watching the Optimizations And Micro-Optimizations In CPython talk by Larry Hastings from PyCon 2010. I never realised that you could create a situation where the nonlocal scope would not be looked in.

Judging performance of python code

Saturday, April 03, 2010

Recently I"ve been messing around with python bytecode, as a result I now know approximately how the python virtual virtual machine works at a bytecode level. I've found this quite interesting but other then satisfying my own curiosity had no benefit from this knowledge. Until today I was wondering which code would be more efficient:

a += 1
if b is not None:
    a += 2
if b is None:
    a += 1
else:
    a += 3

From a style point of view I'd prefer to write the first: I find it slightly more readable since it's got fewer indented blocks and it uses one line less on my screen. But my gut feeling tells me the second is faster, particularly if b is not None (this because in the first sample we'd do 2 add operations instead of one if b is not None). But now I can verify my gut feeling! All I need to do is compile both fragments and use the disassembler to investigate:

>>> co1 = compile(sample1, '', 'exec')
>>> dis.dis(co1)
  2           0 LOAD_NAME                0 (a) 
              3 LOAD_CONST               0 (1) 
              6 INPLACE_ADD          
              7 STORE_NAME               0 (a) 

  3          10 LOAD_NAME                1 (b) 
             13 LOAD_CONST               2 (None) 
             16 COMPARE_OP               9 (is not) 
             19 POP_JUMP_IF_FALSE       35 

  4          22 LOAD_NAME                0 (a) 
             25 LOAD_CONST               1 (2) 
             28 INPLACE_ADD          
             29 STORE_NAME               0 (a) 
             32 JUMP_FORWARD             0 (to 35) 
        >>   35 LOAD_CONST               2 (None) 
             38 RETURN_VALUE
>>> co2 = compile(sample2, '', 'exec')
>>> dis.dis(co2)
  2           0 LOAD_NAME                0 (b) 
              3 LOAD_CONST               2 (None) 
              6 COMPARE_OP               8 (is) 
              9 POP_JUMP_IF_FALSE       25 

  3          12 LOAD_NAME                2 (a) 
             15 LOAD_CONST               0 (1) 
             18 INPLACE_ADD          
             19 STORE_NAME               2 (a) 
             22 JUMP_FORWARD            10 (to 35) 

  5     >>   25 LOAD_NAME                2 (a) 
             28 LOAD_CONST               1 (3) 
             31 INPLACE_ADD          
             32 STORE_NAME               2 (a) 
        >>   35 LOAD_CONST               2 (None) 
             38 RETURN_VALUE

My analysis of this is pretty simple: count the number of instructions for when b is None and for when b is not None.

b is Noneb is not None
sample11015
sample21110

So ultimately the best performance depends on whether b will be None or not. However the difference in the best case is only one instruction, but the difference in the worst case is a whole mighty 4 instructions! This would seem to confirm my gut feeling: the ugly code is better. It also makes me wonder if this is not the sort of optimisation a compiler should be doing: create the bytecode for sample2 regardless of the source code (I'm not a compiler guy and do realise it might not be as simple as that since python is a dynamic language which allows you to change stuff, including executed code, at runtime, yada yada).

There's one more catch tough: I seriously doubt each python instruction can be executed in the same time! So let's use timeit to actually verify this. I'm omitting the trivial code, but this is the result:

$ python3 test.py 
sample1, b is None: 0.128307104111
sample2, b is None: 0.128338098526
sample1, b is not None: 0.244062900543
sample2, b is not None: 0.12109208107

To be honest, the result is exactly as speculated: the first sample is slower when b is not None, all others are pretty much the same. The one odd thing is the pretty much doubling of time for the bad case, this suggest the python virtual machine is spending most of the time doing the INPLACE_ADD instruction while all others are probably very quick.

Anyway, in conclusion I guess you can speculate about performance and get an idea by knowing what bytecode will be generated. But at the end of the day you'll get a simple and better answer by simply using timeit. So knowing something about python bytecode still hasn't gained me any benefit. It was still an interesting exercise tough.

Subscribe to: Posts (Atom)