broker rpc

This commit is contained in:
mahdavi 2024-07-01 17:57:51 +03:30
parent 4b4cd9867b
commit 3d43f956af
6 changed files with 219 additions and 57 deletions

View file

@ -1,3 +1,3 @@
from .celery import app as celery_app
__all__ = ['celery_app']
# from .celery import app as celery_app
#
# __all__ = ['celery_app']

8
accounts/broker_rpc.py Normal file
View file

@ -0,0 +1,8 @@
#!/usr/bin/env python
import os
os.environ.setdefault('DJANGO_SETTINGS_MODULE', 'accounts.settings')
from utils.broker import get_rpc_broker_consumer
channel = get_rpc_broker_consumer()
channel.start_consuming()

View file

@ -1,25 +1,25 @@
from __future__ import absolute_import, unicode_literals
import os
from celery import Celery
from django.conf import settings
# set the default Django settings module for the 'celery' program.
os.environ.setdefault('DJANGO_SETTINGS_MODULE', 'accounts.settings')
app = Celery('notification_service')
# Using a string here means the worker doesn't have to serialize
# the configuration object to child processes.
# - namespace='CELERY' means all celery-related configuration keys
# should have a `CELERY_` prefix.
app.config_from_object('django.conf:settings', namespace='CELERY')
# Load task modules from all registered Django app configs.
app.autodiscover_tasks()
@app.task(bind=True)
def debug_task(self):
print('Request: {0!r}'.format(self.request))
# from __future__ import absolute_import, unicode_literals
#
# import os
#
# from celery import Celery
# from django.conf import settings
#
# # set the default Django settings module for the 'celery' program.
# # os.environ.setdefault('DJANGO_SETTINGS_MODULE', 'accounts.settings')
#
# app = Celery('notification_service')
#
# # Using a string here means the worker doesn't have to serialize
# # the configuration object to child processes.
# # - namespace='CELERY' means all celery-related configuration keys
# # should have a `CELERY_` prefix.
# app.config_from_object('django.conf:settings', namespace='CELERY')
#
# # Load task modules from all registered Django app configs.
# app.autodiscover_tasks()
#
#
# @app.task(bind=True)
# def debug_task(self):
# print('Request: {0!r}'.format(self.request))

View file

@ -13,7 +13,7 @@ from django.utils.translation import gettext_lazy as _
import uuid
from .provinces_and_cities import state
from .tasks import send_notification
# from .tasks import send_notification
from apps.users.constans import MAX_OTP_TRY, DEVELOPMENT_PHONE_NUMBERS

View file

@ -1,28 +1,28 @@
from celery import shared_task
from django.conf import settings
# from service_clients import Client, AccountsClient
# SCOPES = ['notifications.notification:submit']
# client = Client(client_id=settings.CLIENT_ID,
# client_secret=settings.CLIENT_SECRET,
# scopes=SCOPES,
# grant_type=AccountsClient.GRANT_CLIENT_CREDENTIALS)
@shared_task
def send_notification(phone_number=None,
title=None,
body=None,
user_uuid=None,
email=None,
notification_type=1):
data = {
'phone_number': phone_number,
"body": f'{title}:\n{body}'
}
try:
print(client.request(required_scopes=SCOPES, url=f'{settings.BASE_NOTIFICATION_URL}/notifications/', data=data,
method='post', timeout=5).text)
except Exception as e:
print(e)
# from celery import shared_task
# from django.conf import settings
# # from service_clients import Client, AccountsClient
#
# # SCOPES = ['notifications.notification:submit']
# # client = Client(client_id=settings.CLIENT_ID,
# # client_secret=settings.CLIENT_SECRET,
# # scopes=SCOPES,
# # grant_type=AccountsClient.GRANT_CLIENT_CREDENTIALS)
#
# @shared_task
# def send_notification(phone_number=None,
# title=None,
# body=None,
# user_uuid=None,
# email=None,
# notification_type=1):
#
# data = {
# 'phone_number': phone_number,
# "body": f'{title}:\n{body}'
# }
# try:
# print(client.request(required_scopes=SCOPES, url=f'{settings.BASE_NOTIFICATION_URL}/notifications/', data=data,
# method='post', timeout=5).text)
# except Exception as e:
# print(e)
#

154
utils/broker.py Normal file
View file

@ -0,0 +1,154 @@
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)
print(response.data)
print(dir(response))
# print(views.TestApiView().dispatch(request).data)
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=str(response.data),
)
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")
return channel
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)