MCPcopy Create free account
hub / github.com/apache/qpid-proton / Send

Class Send

python/examples/db_send.py:31–96  ·  view source on GitHub ↗

Source from the content-addressed store, hash-verified

29
30
31class 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()

Callers 1

db_send.pyFile · 0.70

Calls

no outgoing calls

Tested by

no test coverage detected