|
1 | 1 | from rabbitmq_amqp_python_client import ( |
2 | 2 | BindingSpecification, |
3 | 3 | Connection, |
| 4 | + Event, |
4 | 5 | ExchangeSpecification, |
5 | 6 | Message, |
| 7 | + MessagingHandler, |
6 | 8 | QuorumQueueSpecification, |
7 | 9 | exchange_address, |
| 10 | + queue_address, |
8 | 11 | ) |
9 | 12 |
|
10 | 13 |
|
| 14 | +class MyMessageHandler(MessagingHandler): |
| 15 | + |
| 16 | + def __init__(self): |
| 17 | + super().__init__() |
| 18 | + |
| 19 | + def on_message(self, event: Event): |
| 20 | + print("received message: " + event.message.body) |
| 21 | + self.accept(event.delivery) |
| 22 | + |
| 23 | + def on_connection_closed(self, event: Event): |
| 24 | + print("connection closed") |
| 25 | + |
| 26 | + def on_connection_cloing(self, event: Event): |
| 27 | + print("connection closed") |
| 28 | + |
| 29 | + def on_link_closed(self, event: Event) -> None: |
| 30 | + print("link closed") |
| 31 | + |
| 32 | + def on_rejected(self, event: Event) -> None: |
| 33 | + print("rejected") |
| 34 | + |
| 35 | + |
11 | 36 | def main() -> None: |
12 | 37 | exchange_name = "test-exchange" |
13 | 38 | queue_name = "example-queue" |
@@ -35,21 +60,35 @@ def main() -> None: |
35 | 60 |
|
36 | 61 | addr = exchange_address(exchange_name, routing_key) |
37 | 62 |
|
| 63 | + addr_queue = queue_address(queue_name) |
| 64 | + |
38 | 65 | print("create a publisher and publish a test message") |
39 | 66 | publisher = connection.publisher(addr) |
40 | 67 |
|
41 | 68 | publisher.publish(Message(body="test")) |
42 | 69 |
|
| 70 | + print("purging the queue") |
| 71 | + messages_purged = management.purge_queue(queue_name) |
| 72 | + |
| 73 | + print("messages purged: " + str(messages_purged)) |
| 74 | + |
| 75 | + for i in range(10): |
| 76 | + publisher.publish(Message(body="test")) |
| 77 | + |
43 | 78 | publisher.close() |
44 | 79 |
|
| 80 | + print("create a consumer and consume the test message") |
| 81 | + |
| 82 | + consumer = connection.consumer(addr_queue, handler=MyMessageHandler()) |
| 83 | + |
45 | 84 | print("unbind") |
46 | 85 | management.unbind(bind_name) |
47 | 86 |
|
48 | | - print("purging the queue") |
49 | | - management.purge_queue(queue_name) |
50 | 87 |
|
| 88 | + |
| 89 | + consumer.close() |
51 | 90 | print("delete queue") |
52 | | - management.delete_queue(queue_name) |
| 91 | + #management.delete_queue(queue_name) |
53 | 92 |
|
54 | 93 | print("delete exchange") |
55 | 94 | management.delete_exchange(exchange_name) |
|
0 commit comments