import json from io import BytesIO import django import dpkt import requests from requests import Request from urllib.parse import urlparse from django.core.handlers.base import BaseHandler from django.core import signals from django.urls import set_script_prefix import pika from django.core.handlers.wsgi import WSGIRequest, get_script_name def encode_request(request:Request): pre = request.prepare() header_list = [] for key, value in pre.headers.items(): header_list.append(f'{key}: {value}') header_str = '\n'.join(header_list) header_bytes = f'{pre.method} {pre.url}\n{header_str}'.encode() if pre.body: return b'\n\n'.join([header_bytes, pre.body]) return header_bytes def decode_request(request_data:bytes): d = dpkt.http.Request(request_data) # print(urlparse(d.uri)) # print(d.method) # print(d.headers) # print('d.data', d.data) # print('d.body', d.body) return d # TODO: is it ok to use DjangoRequestFactory class BrokerRequest(WSGIRequest): def get_port(self): """Return the port number for the request as a string.""" return None def get_host(self): return None class BrokerHandler(BaseHandler): request_class = BrokerRequest def __init__(self, *args, **kwargs): super().__init__(*args, **kwargs) self.load_middleware() def __call__(self, environ): set_script_prefix(get_script_name(environ)) signals.request_started.send(sender=self.__class__, environ=environ) request = self.request_class(environ) response = self.get_response(request) response._handler_class = self.__class__ status = "%d %s" % (response.status_code, response.reason_phrase) response_headers = [ *response.items(), *(("Set-Cookie", c.output(header="")) for c in response.cookies.values()), ] if getattr(response, "file_to_stream", None) is not None and environ.get( "wsgi.file_wrapper" ): # If `wsgi.file_wrapper` is used the WSGI server does not call # .close on the response, but on the file wrapper. Patch it to use # response.close instead which takes care of closing all files. response.file_to_stream.close = response.close response = environ["wsgi.file_wrapper"]( response.file_to_stream, response.block_size ) return response def handle_request(requet_bytes, handler): request_data = decode_request(requet_bytes) # TODO: complete envs envirnment = { "REQUEST_METHOD": request_data.method, "PATH_INFO": urlparse(request_data.uri).path, "CONTENT_LENGTH": request_data.headers.get('Content-Length', 0), "CONTENT_TYPE": request_data.headers.get('Content-Type', 'application/json'), "SERVER_NAME": 'kafka', "SERVER_PORT": '0', # "HTTP_X_FORWARDED_PORT": '0', "wsgi.input": BytesIO(request_data.body) } response = handler(envirnment) # print(response.content) return response def get_rpc_broker_consumer(): django.setup(set_prefix=False) broker_handler = BrokerHandler() connection = pika.BlockingConnection( pika.ConnectionParameters(host="localhost"), ) channel = connection.channel() channel.queue_declare(queue="rpc_queue") def on_request(ch, method, props, body): print("response:...") response = handle_request(body, broker_handler) ch.basic_publish( exchange="", routing_key=props.reply_to, properties=pika.BasicProperties(correlation_id=props.correlation_id), body=response.content, ) ch.basic_ack(delivery_tag=method.delivery_tag) channel.basic_qos(prefetch_count=1) channel.basic_consume(queue="rpc_queue", on_message_callback=on_request) print(" [x] Awaiting RPC requests") 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) #