| 29 | |
| 30 | |
| 31 | class Send(MessagingHandler): |
| 32 | def __init__(self, url, count): |
| 33 | super(Send, self).__init__() |
| 34 | self.url = url |
| 35 | self.delay = 0 |
| 36 | self.sent = 0 |
| 37 | self.confirmed = 0 |
| 38 | self.load_count = 0 |
| 39 | self.records = queue.Queue(maxsize=50) |
| 40 | self.target = count |
| 41 | self.db = Db("src_db", EventInjector()) |
| 42 | |
| 43 | def keep_sending(self): |
| 44 | return self.target == 0 or self.sent < self.target |
| 45 | |
| 46 | def on_start(self, event): |
| 47 | self.container = event.container |
| 48 | self.container.selectable(self.db.injector) |
| 49 | self.sender = self.container.create_sender(self.url) |
| 50 | |
| 51 | def on_records_loaded(self, event): |
| 52 | if self.records.empty(): |
| 53 | if event.subject == self.load_count: |
| 54 | print("Exhausted available data, waiting to recheck...") |
| 55 | # check for new data after 5 seconds |
| 56 | self.container.schedule(5, self) |
| 57 | else: |
| 58 | self.send() |
| 59 | |
| 60 | def request_records(self): |
| 61 | if not self.records.full(): |
| 62 | print("loading records...") |
| 63 | self.load_count += 1 |
| 64 | self.db.load(self.records, event=ApplicationEvent( |
| 65 | "records_loaded", link=self.sender, subject=self.load_count)) |
| 66 | |
| 67 | def on_sendable(self, event): |
| 68 | self.send() |
| 69 | |
| 70 | def send(self): |
| 71 | while self.sender.credit and not self.records.empty(): |
| 72 | if not self.keep_sending(): |
| 73 | return |
| 74 | record = self.records.get(False) |
| 75 | id = record['id'] |
| 76 | self.sender.send(Message(id=id, durable=True, body=record['description']), tag=str(id)) |
| 77 | self.sent += 1 |
| 78 | print("sent message %s" % id) |
| 79 | self.request_records() |
| 80 | |
| 81 | def on_settled(self, event): |
| 82 | id = int(event.delivery.tag) |
| 83 | self.db.delete(id) |
| 84 | print("settled message %s" % id) |
| 85 | self.confirmed += 1 |
| 86 | if self.confirmed == self.target: |
| 87 | event.connection.close() |
| 88 | self.db.close() |