accounts/utils/broker.py
2024-07-02 14:34:59 +03:30

151 lines
4.2 KiB
Python

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)
#