diff --git a/accounts/broker_rpc.py b/accounts/broker_rpc.py index acbe8e5..6b322c4 100644 --- a/accounts/broker_rpc.py +++ b/accounts/broker_rpc.py @@ -5,4 +5,4 @@ os.environ.setdefault('DJANGO_SETTINGS_MODULE', 'accounts.settings') from utils.broker import get_rpc_broker_consumer channel = get_rpc_broker_consumer() -channel.start_consuming() + diff --git a/utils/broker.py b/utils/broker.py index f1dc0a0..99e6bff 100644 --- a/utils/broker.py +++ b/utils/broker.py @@ -104,7 +104,7 @@ def handle_request(requet_bytes, handler): response = handler(envirnment) print(response.content) - print(response.data) + # print(response.data) print(dir(response)) # print(views.TestApiView().dispatch(request).data) return response @@ -130,7 +130,7 @@ def get_rpc_broker_consumer(): exchange="", routing_key=props.reply_to, properties=pika.BasicProperties(correlation_id=props.correlation_id), - body=str(response.data), + body=str(response.content), ) ch.basic_ack(delivery_tag=method.delivery_tag) @@ -138,17 +138,17 @@ def get_rpc_broker_consumer(): channel.basic_consume(queue="rpc_queue", on_message_callback=on_request) print(" [x] Awaiting RPC requests") - return channel + channel.start_consuming() -if __name__ == '__main__': - method = 'GET' - url = 'http://accounts/test_app/test_api?test=1&test=2&test=3' - - request = Request(method=method, url=url, json={'key_1': 'requests_lib'}) - data = encode_request(request) - # print(data) - - request_dict = decode_request(data) - print(request_dict) - +# if __name__ == '__main__': +# method = 'GET' +# url = 'http://accounts/test_app/test_api?test=1&test=2&test=3' +# +# request = Request(method=method, url=url, json={'key_1': 'requests_lib'}) +# data = encode_request(request) +# # print(data) +# +# request_dict = decode_request(data) +# print(request_dict) +# diff --git a/utils/rpc_client.py b/utils/rpc_client.py new file mode 100644 index 0000000..a3976e2 --- /dev/null +++ b/utils/rpc_client.py @@ -0,0 +1,66 @@ +#!/usr/bin/env python +import uuid +import pika +from requests import Request + +from utils.broker import encode_request + + +class RpcClient(object): + def __init__(self): + self.connection = pika.BlockingConnection( + pika.ConnectionParameters(host="localhost"), + ) + + self.channel = self.connection.channel() + + result = self.channel.queue_declare(queue="", exclusive=True) + self.callback_queue = result.method.queue + + self.channel.basic_consume( + queue=self.callback_queue, + on_message_callback=self.on_response, + auto_ack=True, + ) + + self.response = None + self.corr_id = None + + def on_response(self, ch, method, props, body): + if self.corr_id == props.correlation_id: + self.response = body + + def call(self, request, correlation_id=None): + data = encode_request(request) + self.response = None + if not correlation_id: + self.corr_id = str(uuid.uuid4()) + + + + self.channel.basic_publish( + exchange="", + routing_key="rpc_queue", + properties=pika.BasicProperties( + reply_to=self.callback_queue, + correlation_id=self.corr_id, + ), + body=data, + ) + self.connection.process_data_events(time_limit=5) + return self.response + + + +# +method = 'GET' +url = 'http://accounts/test_app/test_api?test=1&test=2&test=3' + +request = Request(method=method, url=url, json={'key_1': 'requests_lib'}) + + +rpc = RpcClient() + +print(" [x] Requesting ...") +response = rpc.call(request) +print(response) \ No newline at end of file