summaryrefslogtreecommitdiff
path: root/joblib.py
blob: 4cc3b94ef0846c0ee6d26b72d9c743a5f96f372c (plain)
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
#!/usr/bin/python3

import json
import logging
import pika


class JobQueue(object):
    def __init__(self):
        self._qcon = None
        self._qchan = None

    def connect(self, **args):
        """Connect to RabbitMQ."""
        logging.info("connecting to broker")
        self._qcon = pika.BlockingConnection(pika.ConnectionParameters(
            **args))
        self._qchan = self._qcon.channel()

    def request(self, queue, payload):
        """Create a new job and return its id."""
        body = json.dumps(payload)
        self._qchan.queue_declare(queue=queue, durable=True)
        self._qchan.basic_publish(exchange="",
                                  routing_key=queue,
                                  body=body,
                                  properties=pika.BasicProperties(
                                      delivery_mode = 2
                                  ))

class JobFailed(Exception):
    pass


class Worker(JobQueue):
    """Abstract service worker class."""
    def __init__(self, jobtype):
        super().__init__()
        self._queue = jobtype
        self._up = False
        self._notifq = []

    def setup(self):
        # override
        pass

    def work(self, payload):
        # override
        pass

    def commit(self, failed):
        # override
        pass

    def request(self, queue, payload, now=True):
        if now:
            super().request(queue, payload)
        else:
            # send after COMMIT
            self._notifq.append((queue, payload))

    def run(self):
        self._qchan.queue_declare(queue=self._queue, durable=True)
        self._qchan.basic_qos(prefetch_count=1)
        for msg in self._qchan.consume(queue=self._queue,
                                       inactivity_timeout=1):
            # TODO force commit after n, c+=1
            if msg is None:
                # idling
                self.commit(False)
                # send notifs dependant on COMMIT
                for notif in self._notifq:
                    self.request(*notif)
                self._notifq = []
                continue

            method, properties, body = msg
            # run a job
            payload = json.loads(body)
            try:
                self.work(payload)
            except JobFailed as err:
                pass  # TODO
            self._qchan.basic_ack(delivery_tag=method.delivery_tag)