diff --git a/accounts/__init__.py b/accounts/__init__.py index cd04264..252c120 100644 --- a/accounts/__init__.py +++ b/accounts/__init__.py @@ -1,3 +1,3 @@ -from .celery import app as celery_app - -__all__ = ['celery_app'] +# from .celery import app as celery_app +# +# __all__ = ['celery_app'] diff --git a/accounts/broker_rpc.py b/accounts/broker_rpc.py new file mode 100644 index 0000000..acbe8e5 --- /dev/null +++ b/accounts/broker_rpc.py @@ -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() diff --git a/accounts/celery.py b/accounts/celery.py index 80ea358..6b7fc88 100644 --- a/accounts/celery.py +++ b/accounts/celery.py @@ -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)) diff --git a/apps/users/models.py b/apps/users/models.py index 189e73c..5da1e0e 100644 --- a/apps/users/models.py +++ b/apps/users/models.py @@ -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 diff --git a/apps/users/tasks.py b/apps/users/tasks.py index 000225c..f11dd9b 100644 --- a/apps/users/tasks.py +++ b/apps/users/tasks.py @@ -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) +# diff --git a/utils/broker.py b/utils/broker.py new file mode 100644 index 0000000..f1dc0a0 --- /dev/null +++ b/utils/broker.py @@ -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) +