From a81fdb6dd4dfe349e76d4f1575394e69182baf7f Mon Sep 17 00:00:00 2001 From: Khushal P Soonderji Date: Tue, 30 Sep 2025 14:08:04 +0530 Subject: [PATCH] 123COMMENT --- .../{socketio/__init__.py => sio_client.py} | 0 playground/socketio_test/__init__.py | 0 .../{socketio => socketio_test}/from_kafka.py | 0 .../kafka_passthrough.py | 0 .../tick_simulator.py | 0 .../tick_simulator_v2.py | 0 .../{socketio => socketio_test}/to_kafka.py | 0 utils_v2/api/response.py | 9 + wsio/finstitutions/trading/tick_out_v2.py | 653 ++++++++++++++++++ wsio_v2/__init__.py | 0 wsio_v2/app.py | 202 ++++++ wsio_v2/app_state.py | 379 ++++++++++ wsio_v2/finstitutions/__init__.py | 0 wsio_v2/finstitutions/trading/__init__.py | 0 wsio_v2/finstitutions/trading/tick_out.py | 139 ++++ wsio_v2/helpers/__init__.py | 0 wsio_v2/helpers/helpers.py | 398 +++++++++++ wsio_v2/test/__init__.py | 0 wsio_v2/test/echo.py | 137 ++++ 19 files changed, 1917 insertions(+) rename playground/{socketio/__init__.py => sio_client.py} (100%) create mode 100644 playground/socketio_test/__init__.py rename playground/{socketio => socketio_test}/from_kafka.py (100%) rename playground/{socketio => socketio_test}/kafka_passthrough.py (100%) rename playground/{socketio => socketio_test}/tick_simulator.py (100%) rename playground/{socketio => socketio_test}/tick_simulator_v2.py (100%) rename playground/{socketio => socketio_test}/to_kafka.py (100%) create mode 100644 wsio/finstitutions/trading/tick_out_v2.py create mode 100644 wsio_v2/__init__.py create mode 100644 wsio_v2/app.py create mode 100644 wsio_v2/app_state.py create mode 100644 wsio_v2/finstitutions/__init__.py create mode 100644 wsio_v2/finstitutions/trading/__init__.py create mode 100644 wsio_v2/finstitutions/trading/tick_out.py create mode 100644 wsio_v2/helpers/__init__.py create mode 100644 wsio_v2/helpers/helpers.py create mode 100644 wsio_v2/test/__init__.py create mode 100644 wsio_v2/test/echo.py diff --git a/playground/socketio/__init__.py b/playground/sio_client.py similarity index 100% rename from playground/socketio/__init__.py rename to playground/sio_client.py diff --git a/playground/socketio_test/__init__.py b/playground/socketio_test/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/playground/socketio/from_kafka.py b/playground/socketio_test/from_kafka.py similarity index 100% rename from playground/socketio/from_kafka.py rename to playground/socketio_test/from_kafka.py diff --git a/playground/socketio/kafka_passthrough.py b/playground/socketio_test/kafka_passthrough.py similarity index 100% rename from playground/socketio/kafka_passthrough.py rename to playground/socketio_test/kafka_passthrough.py diff --git a/playground/socketio/tick_simulator.py b/playground/socketio_test/tick_simulator.py similarity index 100% rename from playground/socketio/tick_simulator.py rename to playground/socketio_test/tick_simulator.py diff --git a/playground/socketio/tick_simulator_v2.py b/playground/socketio_test/tick_simulator_v2.py similarity index 100% rename from playground/socketio/tick_simulator_v2.py rename to playground/socketio_test/tick_simulator_v2.py diff --git a/playground/socketio/to_kafka.py b/playground/socketio_test/to_kafka.py similarity index 100% rename from playground/socketio/to_kafka.py rename to playground/socketio_test/to_kafka.py diff --git a/utils_v2/api/response.py b/utils_v2/api/response.py index 7783cfe..740232c 100644 --- a/utils_v2/api/response.py +++ b/utils_v2/api/response.py @@ -79,6 +79,15 @@ class ResponseModel(BaseModel): http_code: Optional[HttpCodes] = None api_version: Optional[str] = None + @property + def success(self) -> bool: + + """ + A quick wy to check if the response indicates a successful outcome. + """ + + return True if self.status_code.value[0] else False + def for_quart(self): """ diff --git a/wsio/finstitutions/trading/tick_out_v2.py b/wsio/finstitutions/trading/tick_out_v2.py new file mode 100644 index 0000000..823c3dc --- /dev/null +++ b/wsio/finstitutions/trading/tick_out_v2.py @@ -0,0 +1,653 @@ +""" + + AUTHOR: + + Khushal P Soonderji + + DATE: + + Thursday, 2nd Jan., 2025 + + OBJECTIVE: + + To broadcast live tick updates to connected clients. It doesn't matter which stockbroker we are getting the + ticks from as long as we are reading standardized ticks from the Kafka queue. + + REFERENCES: + + N/A + + DOWNLOADS: + + N/A + +""" + + +# ***************************************************************************************************************** +# ***** **** +# *** IMPORT *** +# ***** **** +# ***************************************************************************************************************** + + +# To make sibling directories accessible for imports: +import sys +sys.path.append(".") +sys.path.append("..") + +# System-level activities: +import io +import os +import random + +# My utils: +from utils_v2.string import json +from utils_v2.string import regex +from utils_v2.system import files +from utils_v2.date_time import date_time +from utils_v2.database.async_mongo_v2 import AsyncMongo +from utils_v2.queue.kafka.controllers.async_kafka import ConsumerKafka, get_ssl_context +from utils_v2.cache.async_redis_cache_v3 import AsyncRedisCache +from utils_v2.serialization.json_serializer import JSONSerializer + +# To make HTTP calls: +import httpx + +# To work with date and time: +import datetime +import time + +# Models: +from models.core.user import CoreUserInfoModel + +# To work with SocketIO: +import socket +import socketio + +# For asynchronous activities: +import asyncio + +# To work with various datatypes: +from typing import List + +# Debugging: +from icecream import IceCreamDebugger + + +# ***************************************************************************************************************** +# ***** **** +# *** MACROS / ONE-TIME INIT *** +# ***** **** +# ***************************************************************************************************************** + + +# Debugging: +printer = IceCreamDebugger(prefix = "Tick-Out | ", includeContext = True) +no_context_printer = IceCreamDebugger(prefix = "Tick-Out | ", includeContext = False) + +# To make API calls: +http_client = httpx.AsyncClient( + limits = httpx.Limits( + max_connections = 100, # ............ Maximum number of connections allowed in the pool. + max_keepalive_connections = 50, # ... Maximum number of connections that can be kept alive. + ), + timeout = httpx.Timeout( + pool = 120.0, # .... Time to wait for a free connection from the pool. + connect = 2.5, # ... Time to wait for establishing a connection to the server. + write = 10.0, # .... Time to wait for sending data. + read = 9.9 # ....... Time to wait for receiving data. + ) +) + +# General: +SERVER_HOSTNAME = str(socket.gethostname()) + +# For SocketIO: +# Namespaces: +NAMESPACE_MODULE = None +NAMESPACE_PASSTHROUGH = "/passthrough" +# Events: +EVENT_CONNECT = "connect" +EVENT_DISCONNECT = "disconnect" +EVENT_ECHO = "echo" +EVENT_TICKS = "ticks" + + +# ***************************************************************************************************************** +# ***** **** +# *** VARIABLES *** +# ***** **** +# ***************************************************************************************************************** + + +# For SocketIO: +ALLOWED_ORIGINS = [] +sio = socketio.AsyncServer( + cors_allowed_origins = "*", + async_mode = "asgi" +) +app = socketio.ASGIApp(sio) + +# Redis: +redis_cache: AsyncRedisCache | None = None + +# For kafka: +kafka_consumer: ConsumerKafka | None = None + +# Session-awareness and maintenance of this script's state: +SCRIPT_DATA = {} +exclusive_lock = asyncio.Semaphore(1) +CONNECTED_CLIENTS = {} +FLAGS = { + "initDone": False +} + + +# ***************************************************************************************************************** +# ***** **** +# *** FUNCTIONS *** +# ***** **** +# ***************************************************************************************************************** + + +def origin_is_allowed(origin: str) -> bool: + + """ + To check if a given origin is in the allowed list. + :param origin: The origin of your request. + :return: True if allowed, else False. + """ + + # Start by assuming failure: + is_allowed = False + + # Check through all the allowed origins: + for allowed in ALLOWED_ORIGINS: + try: + if regex.match(origin, allowed): + is_allowed = True + break + except Exception as exception: + printer(exception) + + # Done here: + return is_allowed + + +# --------------------------------------------------------------------------------------------------------------------- + + +async def send_ticks(ticks: List[dict]) -> None: + + """ + Here's where we decide which client gets which tick and send it out. + WARNING: WE ARE ASSUMING THAT NO FURTHER FORMATING/COMPUTATION IS REQUIRED OTHER THAN SELECTING WHICH SUBSETS OF + TICKS TO SEND TO WHICH CLIENTS. FOR US THE TICKS ALREADY HAVE ALL THE DATA NEEDED TO BE SEND TO + RESPECTIVE CLIENTS. + :param ticks: The list of individual tick updates to send out to the clients. + :return: None + """ + + # Currently we're just broadcasting + # all the data to all the clients: + await sio.emit( + event = EVENT_TICKS, + data = ticks, + namespace = NAMESPACE_MODULE + ) + + +# --------------------------------------------------------------------------------------------------------------------- + + +async def ticks_from_kafka( + consumer: ConsumerKafka, + fetch_count: int = 100, + fetch_timeout: float = 1.0 +) -> None: + + """ + This function must run in the background forever and just keep listening for ticks on Kafka and keep relaying them + to all the connected clients as per their watchlists. + :param consumer: The preconfigured Kafka consumer that can listen for ticks in asynchronous mode. + :param fetch_count: How many messages to consume in one go. + :param fetch_timeout: How long to wait (in seconds) while consuming messages from Kafka. + :return: None + """ + + printer("Starting Kafka consumer (ticks).") + + # Do the next part infinitely: + while True: + + # Note the time: + now_utc = date_time.get_current_utc_date_time().timestamp() + + # Get messages form Kafka: + ticks = await consumer.consume( + count = fetch_count, + timeout = fetch_timeout + ) + + # If there are no updates to give: + if not ticks: + no_context_printer("No ticks.") + continue + + # Each message must be treated as an array of tick updates (list of dicts). + # In case the producer is sending each individual tick as a separate message, + # we normalize it to be a list: + tasks = [send_ticks(t.value if isinstance(t.value, list) else [t.value]) for t in ticks] + results = await asyncio.gather(*tasks) + + # Analyze the ticks: + # latency = [abs(now_utc - t.value.get("rcvdTs", t.value["tradeTs"])) for t in ticks] + latency = [abs(now_utc - t.ts.timestamp()) for t in ticks] + avg_latency = sum(latency) / len(latency) + total_ticks = len(ticks) + # late_cutoff_seconds = 3.0 + # + # late_ticks = 0 + # for tick in ticks: + # if now_utc - tick.value["tradeTs"] > late_cutoff_seconds: + # late_ticks += 1 + # ticks_str = f"COUNT: {total_ticks: >5,} | LATE: {late_ticks: >5,} ({(late_ticks/total_ticks)*100.0:.2f}%)" + ticks_str = f"COUNT: {total_ticks: >5,} | AVG. LATENCY: {avg_latency:.5f}" + no_context_printer(ticks_str) + + +# --------------------------------------------------------------------------------------------------------------------- +# DUMMY TICKS WSIO TEST - 23 - 06 - 2025 +# ---------------------------------------------------------------------------------------------------------------------- + +async def dummy_ticks( +) -> None: + + """ + This function must run in the background forever and just keep listening for ticks on Kafka and keep relaying them + to all the connected clients as per their watchlists. + :return: None + """ + + printer("Starting Kafka consumer (ticks).") + + # Do the next part infinitely: + while True: + no_context_printer("DUMMY TICKS") + + # Note the time: + now_utc = date_time.get_current_utc_date_time().timestamp() + + # Get messages form Kafka: + ticks = [ + { + "broker": "zerodhaKite", + "brokerToken": 999999, + "exchange": "NSE", + "exchangeToken": "999999", + "segment": "NSE", + "type": "EQ", + "symbol": "TCAOFF", + "name": "The CA Office", + "expiry": None, + "strike": 0, + "bidQty": 65, + "bidRate": 1715, + "askQty": 105, + "askRate": 1715.5, + "o": 1675, + "h": 1723.6, + "l": 1665, + "ltp": 1715.05*random.uniform(0.95, 1.05), + "qty": 20, + "chg": 29.200000000000045, + "pChg": 1.732115316170367, + "vwap": 1700.16, + "totVol": 4810200, + "rcvdTs": now_utc, + "tradeTs": now_utc, + "tradeTz": "Asia/Kolkata", + "exchgTs": now_utc, + "exchgTz": "Asia/Kolkata" + } + ] + + # If there are no updates to give: + if not ticks: continue + + # Each message must be treated as an array of tick updates (list of dicts). + # In case the producer is sending each individual tick as a separate message, + # we normalize it to be a list: + tasks = [send_ticks(ticks)] + results = await asyncio.gather(*tasks) + + await asyncio.sleep(0.25) + + # Analyze the ticks: + # latency = [abs(now_utc - t.value.get("rcvdTs", t.value["tradeTs"])) for t in ticks] + # latency = [abs(now_utc - t.ts.timestamp()) for t in ticks] + # avg_latency = sum(latency) / len(latency) + # total_ticks = len(ticks) + # late_cutoff_seconds = 3.0 + # + # late_ticks = 0 + # for tick in ticks: + # if now_utc - tick.value["tradeTs"] > late_cutoff_seconds: + # late_ticks += 1 + # ticks_str = f"COUNT: {total_ticks: >5,} | LATE: {late_ticks: >5,} ({(late_ticks/total_ticks)*100.0:.2f}%)" + # ticks_str = f"COUNT: {total_ticks: >5,} | AVG. LATENCY: {avg_latency:.5f}" + # no_context_printer(ticks_str) + + + +# --------------------------------------------------------------------------------------------------------------------- + + +async def init( + script_id: str, + debug: bool +): + + """ + To initialize all credentials, instances, and connectivity for this whole script. + :param script_id: The id to use to load cred and data from the internal service. + :param debug: Whether, or not, you would like to print the debug messages. + :return: True if initialized successfully, else False. + """ + + # Declare the required global variables: + global SCRIPT_DATA + global ALLOWED_ORIGINS + global redis_cache + global kafka_consumer + + # Basic stuff: + if debug: printer.enable() + printer("Initializing.") + + # ┏┓ • • + # ┃┃┏┓┓┏┓┓┏┓┏ + # ┗┛┛ ┗┗┫┗┛┗┛ + # ┛ + + response = await http_client.post( + url = r"https://api.thecaoffice.com/ca/get/title", + headers = {"Origin": "https://thecaoffice.com/"}, + data = { + "domainName": "127.0.0.1:1234", + "screenWidth": 1920, + "screenHeight": 1080 + } + ) + if response.status_code not in [200]: + print("FATAL: ALLOWED ORIGINS NOT FETCHED!") + return False + ALLOWED_ORIGINS = [origin["domain"] for origin in response.json().get("data", {}).get("rs2", [])] + printer(ALLOWED_ORIGINS) + if len(ALLOWED_ORIGINS) < 1: + print("FATAL: ALLOWED ORIGINS IS EMPTY!") + return False + + # ┏┓ ┓ ┓ ┳┓ + # ┃ ┏┓┏┓┏┫ ┏┓┏┓┏┫ ┃┃┏┓╋┏┓ + # ┗┛┛ ┗ ┗┻ ┗┻┛┗┗┻ ┻┛┗┻┗┗┻ + + # Get the script credentials: + response = await http_client.get( + url = r"https://nexcom.ditscentre.in/internal/cred/get", + headers = {"X-Script-Id": script_id} + ) + if response.status_code not in [200]: + print("FATAL: SCRIPT CREDENTIALS LOADING FAILED!") + return False + script_cred = response.json().get("data") + + # Get the script data: + response = await http_client.get( + url = r"https://nexcom.ditscentre.in/internal/data/get", + headers = {"X-Script-Id": script_id} + ) + if response.status_code not in [200]: + print("FATAL: SCRIPT DATA LOADING FAILED!") + return False + SCRIPT_DATA = response.json().get("data") + + # Done with this step: + printer("Cred and Data loaded.") + + # ┓┏┓ ┏┓ ┏┓┓• + # ┃┫ ┏┓╋┃┏┏┓ ┃ ┃┓┏┓┏┓╋┏ + # ┛┗┛┗┻┛┛┗┗┻ ┗┛┗┗┗ ┛┗┗┛ + + # Create the consumer that will listen to changes in watchlist: + consumer_creds = script_cred["kafka"]["consumer"] + kafka_consumer = ConsumerKafka( + topic = consumer_creds["topic"], + bootstrap_servers = consumer_creds["config"]["bootstrapServers"], + security_protocol = consumer_creds["config"].get("securityProtocol", "PLAINTEXT"), + ssl_context = get_ssl_context( + ca_file = consumer_creds["config"].get("caFile"), + cert_file = consumer_creds["config"].get("certFile"), + key_file = consumer_creds["config"].get("keyFile"), + ), + serializer = JSONSerializer(), + debug = debug + ) + if not await kafka_consumer.connect(): + print("FATAL: KAFKA CONSUMER NOT CREATED!") + return False + printer("Kafka consumer ready.") + + # ┳┓ ┓• ┏┓ ┓ + # ┣┫┏┓┏┫┓┏ ━━ ┃ ┏┓┏┣┓┏┓ + # ┛┗┗ ┗┻┗┛ ┗┛┗┻┗┛┗┗ + + redis_cache = AsyncRedisCache( + connection_string=script_cred["redisCache"]["general"]["sentinelJson"], + # connection_string = script_cred["redisCache"]["general"]["connectionString"], + # serializer = JSONSerializer(), + debug = debug, + debug_prefix = "General Cache | " + ) + if not await redis_cache.connect(): + print("FATAL: REDIS CACHE NOT CREATED!") + return False + printer("Redis cache ready.") + + # ┳┓ ┓ ┓ ┏┳┓ ┓ + # ┣┫┏┓┏┃┏┏┓┏┓┏┓┓┏┏┓┏┫ ┃ ┏┓┏┃┏┏ + # ┻┛┗┻┗┛┗┗┫┛ ┗┛┗┻┛┗┗┻ ┻ ┗┻┛┛┗┛ + # ┛ + + # Start the background task that will receive ticks from the Kafka queue and broadcast them to the respective + # connected clients: + sio.start_background_task( + ticks_from_kafka, + consumer = kafka_consumer, + fetch_count = 1_000, + fetch_timeout = 1.0 + ) + + sio.start_background_task( + dummy_ticks + ) + + # ┳┓ + # ┃┃┏┓┏┓┏┓ + # ┻┛┗┛┛┗┗ + + # If everything went well, we return with success: + printer("Initialization done.") + return True + + +# --------------------------------------------------------------------------------------------------------------------- + + +@sio.on(event = EVENT_CONNECT, namespace = NAMESPACE_MODULE) +async def on_connect(sid, environ, *args) -> bool: + + """ + The event handler for when a new connection request comes in. + :param sid: The session id of the incoming request (generated by SocketIO). + :param environ: The set of headers and other connection-specific values. + :param args: Any extra input coming from the connection request. + :return: True to accept a connection request, False to reject it. + """ + + # declare the required global variables: + global CONNECTED_CLIENTS + + # Initialize the script if needed: + async with exclusive_lock: + if not FLAGS.get("initDone"): + FLAGS["initDone"] = await init( + script_id = os.environ["SCRIPT_ID"], + debug = True if os.environ["DEBUG"].lower() == "true" else False + ) + + # If the initialization failed, we cannot accept the incoming request: + if not FLAGS.get("initDone"): + printer("SOCKET REJECTED: Init. pending.", sid) + return False + + # Check the origin of the incoming request: + printer("Checking origin.") + origin = environ.get("HTTP_ORIGIN", "???") + if not origin_is_allowed(origin): + printer("SOCKET REJECTED: Bad origin.", sid, origin) + return False + + # Get the session token from the incoming request: + printer("Checking session token.") + session_token = environ.get("HTTP_X_SESSION_TOKEN") + if not session_token and len(args) > 0: session_token = args[0].get("X-Session-Token") + if not session_token: + printer("SOCKET REJECTED: No session token.", sid) + return False + + # Get the user's details from the session token: + printer("Fetching user info.") + user_info = await redis_cache.get(key = session_token) + if not user_info: + printer("SOCKET REJECTED: Invalid session token.", sid) + return False + user_info = CoreUserInfoModel(**user_info) + + # Get the user's watchlist and note down the details. + # Consider the following structure for a user's info: + redis_key = f"io_{session_token}" + async with exclusive_lock: + CONNECTED_CLIENTS[sid] = { + "user": user_info, + "redisKey": redis_key, + "rooms": [] + } + sid_cached = await redis_cache.set( + key = redis_key, + value = {"server": SERVER_HOSTNAME, "socket_id": sid} + ) + + # Done here: + printer("SOCKET ACCEPTED.", sid, sid_cached) + return True + + +# --------------------------------------------------------------------------------------------------------------------- + + +@sio.on(event = EVENT_DISCONNECT, namespace = NAMESPACE_MODULE) +async def handle_disconnect(sid, reason) -> None: + + """ + To handle a disconnect event. Automatically triggered when a client disconnects from the server. + :param sid: The session id of the client (generated by SocketIO on connecting). + :param reason: The hint about why the disconnection happened. + :return: None. + """ + + # declare the required global variables: + global CONNECTED_CLIENTS + + # register the disconnect in the global variable, and on the cache server: + client_info = {} + async with exclusive_lock: client_info = CONNECTED_CLIENTS.pop(sid, None) + sid_uncached = await redis_cache.delete(key = client_info["redisKey"]) if client_info else False + printer("SOCKET DISCONNECTED", sid, reason, sid_uncached) + + +# ***************************************************************************************************************** +# ***** **** +# *** MAIN PROGRAM *** +# ***** **** +# ***************************************************************************************************************** + + +if __name__ == "__main__": + + printer("Main.") + + # To get args from the terminal: + import argparse + + # To run the ASGI: + import uvicorn + from multiprocessing import freeze_support + + # Get the config from the command-line: + parser = argparse.ArgumentParser(description = f"SocketIO to serve live market data (and a general passthrough).") + parser.add_argument( + "-w", "--workers", + type = int, + help = "The no. of threads to spin up for this instance!", + default = 2 + ) + parser.add_argument( + "-a", "--host", + type = str, + help = "The host for the app. e.g.: '0.0.0.0' or '127.0.0.1'.", + default = "0.0.0.0" + ) + parser.add_argument( + "-p", "--port", + type = int, + help = "The port no. to bind the app to.", + default = 8080 + ) + parser.add_argument( + "-s", "--script-id", + type = str, + help = "The id of this script (will affect the loaded config)." + ) + parser.add_argument( + "-d", "--debug", + action = "store_true", + help = "Whether, or not, you want to see debugging messages in the terminal.", + default = False + ) + args = parser.parse_args() + + # Note down the config; + os.environ["SCRIPT_ID"] = args.script_id + os.environ["DEBUG"] = str(args.debug) + + # Startup message: + printer.enable() + printer(str(args.debug)) + printer.disable() + + # asyncio.run(init( + # script_id=args.script_id, + # debug=args.debug + # )) + + # Run the gateway: + freeze_support() + uvicorn.run( + app = "tick_out:app", + workers = args.workers, + host = args.host, + port = args.port + ) + diff --git a/wsio_v2/__init__.py b/wsio_v2/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/wsio_v2/app.py b/wsio_v2/app.py new file mode 100644 index 0000000..67d39f7 --- /dev/null +++ b/wsio_v2/app.py @@ -0,0 +1,202 @@ +""" + + AUTHOR: + + Khushal P Soonderji + + DATE: + + Create: Saturday, 18th May, 2022 + Update: Thursday, 22nd Aug. 2024 + + OBJECTIVE: + + To provide an easy way to work with '.json' data and files. + + REFERENCES: + + 1) https://www.w3schools.com/python/python_json.asp + + DOWNLOADS: + + N/A + +""" + + +# ***************************************************************************************************************** +# ***** **** +# *** IMPORT *** +# ***** **** +# ***************************************************************************************************************** + + +# To make sibling directories accessible for imports: +import sys +sys.path.append(".") +sys.path.append("..") + +# System-level activities: +import io + +# To work with the JSON standard: +import json + +# To work with files: +from utils_v2.system import files + + +# ***************************************************************************************************************** +# ***** **** +# *** MACROS / ONE-TIME INIT *** +# ***** **** +# ***************************************************************************************************************** + + +# --- Nothing Yet + + +# ***************************************************************************************************************** +# ***** **** +# *** VARIABLES *** +# ***** **** +# ***************************************************************************************************************** + + +# --- Nothing Yet + + +# ***************************************************************************************************************** +# ***** **** +# *** FUNCTIONS *** +# ***** **** +# ***************************************************************************************************************** + + +def from_string(json_data): + + """ + Decodes a JSON string to a pythonic variable like a dict. + :param json_data: The JSON string to decode. + :return: The decoded pythonic variable. + """ + + python_data = json.loads(json_data) + return python_data + + +# --------------------------------------------------------------------------------------------------------------------- + + +def to_string( + python_data, + indent = 4, + default = None, + separators = None, + no_space = False +): + + """ + Converts the given pythonic data to a JSON string. + :param python_data: The input data like a dict. + :param indent: The tab-width for pretty presentation. + :param default: The function to use on something that cannot be directly parsed into a JSON string. + :param separators: Custom separators to use. + :param no_space: If you want a dense JSON string that saves memory by not using spaces or tabs or line-breaks. Not + good for human readability, very good for saving memory. WARNING: THIS OVERRIDES EVERY OTHER PARAMETER EXCEPT + 'default'. + :return: The JSON string representation of the input pythonic data. + """ + + if no_space: + json_data = json.dumps( + python_data, + default = default, + separators = (',', ':') + ) + + else: + json_data = json.dumps( + python_data, + indent = indent, + default = default, + separators = separators + ) + + return json_data + + +# --------------------------------------------------------------------------------------------------------------------- + + +def from_file(file): + + """ + Reads a JSON file and returns it as a pythonic variable like a dict. + :param file: The path to the file on the disk or a file held in RAM as a BytesIO object. + :return: The decoded pythonic variable. + """ + + if isinstance(file, io.BytesIO): + file.seek(0) + json_data = file.getvalue() + else: json_data = files.read_file(file) + python_data = from_string(json_data) + return python_data + + +# --------------------------------------------------------------------------------------------------------------------- + + +def to_file( + file, + python_data, + indent = 4, + default = None, + separators = None, + no_space = False +): + + """ + + :param file: Either a path to a file on disk, or a buffer in RAM in the form of a BytesIO object. + :param python_data: The pythonic data to be converted to the JSON string. + :param indent: The tab-width for pretty presentation. + :param default: The function to use on something that cannot be directly parsed into a JSON string. + :param separators: Custom separators to use. + :param no_space: If you want a dense JSON string that saves memory by not using spaces or tabs or line-breaks. Not + good for human readability, very good for saving memory. WARNING: THIS OVERRIDES EVERY OTHER PARAMETER EXCEPT + 'default'. + :return: True/False if a path was given, else the same BytesIO object with the written JSON data. + """ + + json_data = to_string( + python_data, + indent = indent, + default = default, + separators = separators, + no_space = no_space + ) + + if isinstance(file, io.BytesIO): + file.write(json_data.encode("utf-8")) + file.seek(0) + return file + + else: + try: + files.write_file(file, json_data, mode = "w") + return True + except: return False + + +# ***************************************************************************************************************** +# ***** **** +# *** MAIN PROGRAM *** +# ***** **** +# ***************************************************************************************************************** + + +if __name__ == "__main__": + + pass diff --git a/wsio_v2/app_state.py b/wsio_v2/app_state.py new file mode 100644 index 0000000..24b8d61 --- /dev/null +++ b/wsio_v2/app_state.py @@ -0,0 +1,379 @@ +""" + + AUTHOR: + + Khushal P Soonderji + + DATE: + + Create: Monday, 29th Sept., 2025 + + OBJECTIVE: + + To have a centralized Socket.IO app from where several namespaces can be registered. This is kind of like how + you can have one Quart app and register several blueprints. + + REFERENCES: + + N/A + + DOWNLOADS: + + N/A + +""" + + +# ***************************************************************************************************************** +# ***** **** +# *** IMPORT *** +# ***** **** +# ***************************************************************************************************************** + + +# To make sibling directories accessible for imports: +import sys +sys.path.append("../wsio") +sys.path.append("..") + +# System-level activities: +import io +import os +import random + +# My utils: +from utils_v2.string import json +from utils_v2.string import regex +from utils_v2.system import files +from utils_v2.date_time import date_time +from utils_v2.database.async_mongo_v2 import AsyncMongo +from utils_v2.queue.kafka.controllers.async_kafka import ConsumerKafka, get_ssl_context +from utils_v2.cache.async_redis_cache_v3 import AsyncRedisCache +from utils_v2.serialization.json_serializer import JSONSerializer + +# To make HTTP calls: +import httpx + +# To work with date and time: +import datetime +import time + +# Models: +from models.core.user import CoreUserInfoModel +from models.finstitutions.trading.ticks import TradingTick + +# To work with SocketIO: +import socket +import socketio + +# For asynchronous activities: +import asyncio + +# To work with various datatypes: +from typing import List + +# Debugging: +from icecream import IceCreamDebugger + + +# ***************************************************************************************************************** +# ***** **** +# *** MACROS / ONE-TIME INIT *** +# ***** **** +# ***************************************************************************************************************** + + +# Debugging: +printer = IceCreamDebugger(prefix = "WSIO | ", includeContext = True) +no_context_printer = IceCreamDebugger(prefix = "WSIO | ", includeContext = False) + +# To make API calls: +http_client = httpx.AsyncClient( + limits = httpx.Limits( + max_connections = 100, # ............ Maximum number of connections allowed in the pool. + max_keepalive_connections = 50, # ... Maximum number of connections that can be kept alive. + ), + timeout = httpx.Timeout( + pool = 120.0, # .... Time to wait for a free connection from the pool. + connect = 2.5, # ... Time to wait for establishing a connection to the server. + write = 10.0, # .... Time to wait for sending data. + read = 9.9 # ....... Time to wait for receiving data. + ) +) + +# General: +SERVER_HOSTNAME = str(socket.gethostname()) + +# For SocketIO: +# Namespaces: +NAMESPACE_TICKS = "/ticks" +NAMESPACE_ORDER_UPDATES = "/order-updates" +# Events: +EVENT_CONNECT = "connect" +EVENT_DISCONNECT = "disconnect" +EVENT_ECHO = "echo" +EVENT_TICKS = "ticks" + + +# ***************************************************************************************************************** +# ***** **** +# *** VARIABLES *** +# ***** **** +# ***************************************************************************************************************** + + +# For SocketIO: +ALLOWED_ORIGINS = [] +sio = socketio.AsyncServer( + cors_allowed_origins = "*", + async_mode = "asgi" +) +app = socketio.ASGIApp(sio) + +# A custom class to maintain the app's state: +class AppState: + def __init__(self): + self.script_data = {} + self.init_done = False + self.connected_clients = {} + self.exclusive_lock = asyncio.Semaphore(1) +app_state = AppState() + + +# ***************************************************************************************************************** +# ***** **** +# *** FUNCTIONS *** +# ***** **** +# ***************************************************************************************************************** + + +def origin_is_allowed(origin: str) -> bool: + + """ + To check if a given origin is in the allowed list. + :param origin: The origin of your request. + :return: True if allowed, else False. + """ + + # Start by assuming failure: + is_allowed = False + + # Check through all the allowed origins: + for allowed in ALLOWED_ORIGINS: + try: + if regex.match(origin, allowed): + is_allowed = True + break + except Exception as exception: + printer(exception) + + # Done here: + return is_allowed + + +# --------------------------------------------------------------------------------------------------------------------- + + +async def init( + script_id: str, + debug: bool +): + + """ + To initialize all credentials, instances, and connectivity for this whole script. + :param script_id: The id to use to load cred and data from the internal service. + :param debug: Whether, or not, you would like to print the debug messages. + :return: True if initialized successfully, else False. + """ + + # Declare the required global variables: + global SCRIPT_DATA + global ALLOWED_ORIGINS + global redis_cache + global kafka_consumer + + # Basic stuff: + if debug: printer.enable() + printer("Initializing.") + + # ┏┓ • • + # ┃┃┏┓┓┏┓┓┏┓┏ + # ┗┛┛ ┗┗┫┗┛┗┛ + # ┛ + + response = await http_client.post( + url = r"https://api.thecaoffice.com/ca/get/title", + headers = {"Origin": "https://thecaoffice.com/"}, + data = { + "domainName": "127.0.0.1:1234", + "screenWidth": 1920, + "screenHeight": 1080 + } + ) + if response.status_code not in [200]: + print("FATAL: ALLOWED ORIGINS NOT FETCHED!") + return False + ALLOWED_ORIGINS = [origin["domain"] for origin in response.json().get("data", {}).get("rs2", [])] + printer(ALLOWED_ORIGINS) + if len(ALLOWED_ORIGINS) < 1: + print("FATAL: ALLOWED ORIGINS IS EMPTY!") + return False + + # ┏┓ ┓ ┓ ┳┓ + # ┃ ┏┓┏┓┏┫ ┏┓┏┓┏┫ ┃┃┏┓╋┏┓ + # ┗┛┛ ┗ ┗┻ ┗┻┛┗┗┻ ┻┛┗┻┗┗┻ + + # Get the script credentials: + response = await http_client.get( + url = r"https://nexcom.ditscentre.in/internal/cred/get", + headers = {"X-Script-Id": script_id} + ) + if response.status_code not in [200]: + print("FATAL: SCRIPT CREDENTIALS LOADING FAILED!") + return False + script_cred = response.json().get("data") + + # Get the script data: + response = await http_client.get( + url = r"https://nexcom.ditscentre.in/internal/data/get", + headers = {"X-Script-Id": script_id} + ) + if response.status_code not in [200]: + print("FATAL: SCRIPT DATA LOADING FAILED!") + return False + SCRIPT_DATA = response.json().get("data") + + # Done with this step: + printer("Cred and Data loaded.") + + # ┓┏┓ ┏┓ ┏┓┓• + # ┃┫ ┏┓╋┃┏┏┓ ┃ ┃┓┏┓┏┓╋┏ + # ┛┗┛┗┻┛┛┗┗┻ ┗┛┗┗┗ ┛┗┗┛ + + # Create the consumer that will listen to changes in watchlist: + consumer_creds = script_cred["kafka"]["consumer"] + kafka_consumer = ConsumerKafka( + topic = consumer_creds["topic"], + bootstrap_servers = consumer_creds["config"]["bootstrapServers"], + security_protocol = consumer_creds["config"].get("securityProtocol", "PLAINTEXT"), + ssl_context = get_ssl_context( + ca_file = consumer_creds["config"].get("caFile"), + cert_file = consumer_creds["config"].get("certFile"), + key_file = consumer_creds["config"].get("keyFile"), + ), + serializer = JSONSerializer(), + debug = debug + ) + if not await kafka_consumer.connect(): + print("FATAL: KAFKA CONSUMER NOT CREATED!") + return False + printer("Kafka consumer ready.") + + # ┳┓ ┓• ┏┓ ┓ + # ┣┫┏┓┏┫┓┏ ━━ ┃ ┏┓┏┣┓┏┓ + # ┛┗┗ ┗┻┗┛ ┗┛┗┻┗┛┗┗ + + redis_cache = AsyncRedisCache( + connection_string = script_cred["redisCache"]["general"]["sentinelJson"], + # connection_string = script_cred["redisCache"]["general"]["connectionString"], + # serializer = JSONSerializer(), + debug = debug, + debug_prefix = "General Cache | " + ) + if not await redis_cache.connect(): + print("FATAL: REDIS CACHE NOT CREATED!") + return False + printer("Redis cache ready.") + + # ┳┓ + # ┃┃┏┓┏┓┏┓ + # ┻┛┗┛┛┗┗ + + # If everything went well, we return with success: + printer("Initialization done.") + return True + + +# --------------------------------------------------------------------------------------------------------------------- + + +def register_namespaces(): + + """ + A quick function that registers all the namespaces to the same Socket.IO app. + """ + + printer("Namespaces registered.") + + +# ***************************************************************************************************************** +# ***** **** +# *** MAIN PROGRAM *** +# ***** **** +# ***************************************************************************************************************** + + +if __name__ == "__main__": + + printer("Main.") + + # To get args from the terminal: + import argparse + + # To run the ASGI: + import uvicorn + from multiprocessing import freeze_support + + # Get the config from the command-line: + parser = argparse.ArgumentParser(description = f"SocketIO to serve live market data (and a general passthrough).") + parser.add_argument( + "-w", "--workers", + type = int, + help = "The no. of threads to spin up for this instance!", + default = 2 + ) + parser.add_argument( + "-a", "--host", + type = str, + help = "The host for the app. e.g.: '0.0.0.0' or '127.0.0.1'.", + default = "0.0.0.0" + ) + parser.add_argument( + "-p", "--port", + type = int, + help = "The port no. to bind the app to.", + default = 8080 + ) + parser.add_argument( + "-s", "--script-id", + type = str, + help = "The id of this script (will affect the loaded config)." + ) + parser.add_argument( + "-d", "--debug", + action = "store_true", + help = "Whether, or not, you want to see debugging messages in the terminal.", + default = False + ) + args = parser.parse_args() + + # Note down the config; + os.environ["SCRIPT_ID"] = args.script_id + os.environ["DEBUG"] = str(args.debug) + + # Startup message: + printer.enable() + printer(str(args.debug)) + printer.disable() + + # Register all namespaces: + register_namespaces() + + # Run the gateway: + freeze_support() + uvicorn.run( + app = "tick_out:app", + workers = args.workers, + host = args.host, + port = args.port + ) diff --git a/wsio_v2/finstitutions/__init__.py b/wsio_v2/finstitutions/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/wsio_v2/finstitutions/trading/__init__.py b/wsio_v2/finstitutions/trading/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/wsio_v2/finstitutions/trading/tick_out.py b/wsio_v2/finstitutions/trading/tick_out.py new file mode 100644 index 0000000..6db7486 --- /dev/null +++ b/wsio_v2/finstitutions/trading/tick_out.py @@ -0,0 +1,139 @@ +""" + + AUTHOR: + + Khushal P Soonderji + + DATE: + + Create: Tuesday, 30th Sept., 2025 + + OBJECTIVE: + + A simple namespace to test the Socket.IO app with a simple echo utility. + + REFERENCES: + + N/A + + DOWNLOADS: + + N/A + +""" + + +# ***************************************************************************************************************** +# ***** **** +# *** IMPORT *** +# ***** **** +# ***************************************************************************************************************** + + +# To make sibling directories accessible for imports: +import sys +sys.path.append("../wsio") +sys.path.append("..") + +# System-level activities: +import io +import os +import random + +# My utils: +from utils_v2.string import json +from utils_v2.string import regex +from utils_v2.system import files +from utils_v2.date_time import date_time +from utils_v2.database.async_mongo_v2 import AsyncMongo +from utils_v2.queue.kafka.controllers.async_kafka import ConsumerKafka, get_ssl_context +from utils_v2.cache.async_redis_cache_v3 import AsyncRedisCache +from utils_v2.serialization.json_serializer import JSONSerializer + +# To make HTTP calls: +import httpx + +# To work with date and time: +import datetime +import time + +# Models: +from models.core.user import CoreUserInfoModel +from models.finstitutions.trading.ticks import TradingTick + +# To work with SocketIO: +import socket +import socketio + +# To maintain the app's state: +from wsio_v2.app_state import AppState + +# For asynchronous activities: +import asyncio + +# To work with various datatypes: +from typing import List + +# Debugging: +from icecream import IceCreamDebugger + + +# ***************************************************************************************************************** +# ***** **** +# *** MACROS / ONE-TIME INIT *** +# ***** **** +# ***************************************************************************************************************** + + +# --- Nothing Yet + + +# ***************************************************************************************************************** +# ***** **** +# *** VARIABLES *** +# ***** **** +# ***************************************************************************************************************** + + +# --- Nothing Yet + + +# ***************************************************************************************************************** +# ***** **** +# *** FUNCTIONS *** +# ***** **** +# ***************************************************************************************************************** + + +class EchoNamespace(socketio.AsyncNamespace): + + def __init__( + self, + namespace: str, + app_state: AppState, + ): + super().__init__(namespace) + self.app_state = app_state + self.app_state.printer("Registered!") + + async def on_connect(self, sid, environ, *args): + self.app_state.printer("On Connect", sid) + + async def on_disconnect(self, sid, reason, *args): + self.app_state.printer("On Disconnect", sid) + + async def on_echo(self, sid, data): + self.app_state.printer("Event", sid, data, type(data).__name__) + await self.emit("echo", data, to = sid) + + +# ***************************************************************************************************************** +# ***** **** +# *** MAIN PROGRAM *** +# ***** **** +# ***************************************************************************************************************** + + +if __name__ == "__main__": + + pass diff --git a/wsio_v2/helpers/__init__.py b/wsio_v2/helpers/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/wsio_v2/helpers/helpers.py b/wsio_v2/helpers/helpers.py new file mode 100644 index 0000000..8899d7d --- /dev/null +++ b/wsio_v2/helpers/helpers.py @@ -0,0 +1,398 @@ +""" + + AUTHOR: + + Khushal P Soonderji + + DATE: + + Create: Monday, 29th Sept., 2025 + + OBJECTIVE: + + To have a centralized Socket.IO app from where several namespaces can be registered. This is kind of like how + you can have one Quart app and register several blueprints. + + REFERENCES: + + N/A + + DOWNLOADS: + + N/A + +""" + + +# ***************************************************************************************************************** +# ***** **** +# *** IMPORT *** +# ***** **** +# ***************************************************************************************************************** + + +# To make sibling directories accessible for imports: +import sys +sys.path.append(".") +sys.path.append("..") + +# System-level activities: +import io +import os +import random + +# My utils: +from utils_v2.string import json +from utils_v2.string import regex +from utils_v2.system import files +from utils_v2.date_time import date_time +from utils_v2.database.async_mongo_v2 import AsyncMongo +from utils_v2.queue.kafka.controllers.async_kafka import ConsumerKafka, get_ssl_context +from utils_v2.cache.async_redis_cache_v3 import AsyncRedisCache +from utils_v2.serialization.json_serializer import JSONSerializer + +# To make HTTP calls: +import httpx + +# To work with date and time: +import datetime +import time + +# Models: +from models.core.user import CoreUserInfoModel +from models.finstitutions.trading.ticks import TradingTick + +# To work with SocketIO: +import socket +import socketio + +# To maintain the app's state: +from wsio_v2.app_state import AppState +from wsio_v2.test.test_ns import TestNamespace +from wsio_v2.test.echo import echo + +# For asynchronous activities: +import asyncio + +# To work with various datatypes: +from typing import List + +# Debugging: +from icecream import IceCreamDebugger + + +# ***************************************************************************************************************** +# ***** **** +# *** MACROS / ONE-TIME INIT *** +# ***** **** +# ***************************************************************************************************************** + + +# A custom class to maintain the app's state: +app_state = AppState() + +# Debugging: +app_state.printer = IceCreamDebugger(prefix = "WSIO | ", includeContext = True) +app_state.no_context_printer = IceCreamDebugger(prefix = "WSIO | ", includeContext = False) + +# To make API calls: +app_state.http_client = httpx.AsyncClient( + limits = httpx.Limits( + max_connections = 100, # ............ Maximum number of connections allowed in the pool. + max_keepalive_connections = 50, # ... Maximum number of connections that can be kept alive. + ), + timeout = httpx.Timeout( + pool = 120.0, # .... Time to wait for a free connection from the pool. + connect = 2.5, # ... Time to wait for establishing a connection to the server. + write = 10.0, # .... Time to wait for sending data. + read = 9.9 # ....... Time to wait for receiving data. + ) +) + +# General: +app_state.SERVER_HOSTNAME = str(socket.gethostname()) +app_state.ALLOWED_ORIGINS = [] + +# Namespaces: +app_state.NAMESPACE_DEFAULT = "/" +app_state.NAMESPACE_TEST = "/test" +app_state.NAMESPACE_TICKS = "/ticks" +app_state.NAMESPACE_ORDER_UPDATES = "/order-updates" + +# Events: +app_state.EVENT_CONNECT = "connect" +app_state.EVENT_DISCONNECT = "disconnect" +app_state.EVENT_ECHO = "echo" +app_state.EVENT_TICKS = "ticks" + + +# ***************************************************************************************************************** +# ***** **** +# *** VARIABLES *** +# ***** **** +# ***************************************************************************************************************** + + +# For SocketIO: +sio = socketio.AsyncServer( + cors_allowed_origins = "*", + async_mode = "asgi" +) +app = socketio.ASGIApp(sio) +sio.register_namespace(TestNamespace("/test")) + + +# ***************************************************************************************************************** +# ***** **** +# *** FUNCTIONS *** +# ***** **** +# ***************************************************************************************************************** + + +def origin_is_allowed(origin: str) -> bool: + + """ + To check if a given origin is in the allowed list. + :param origin: The origin of your request. + :return: True if allowed, else False. + """ + + # For debugging: + app_state.printer("Validating origin.") + + # Start by assuming failure: + is_allowed = False + + # Check through all the allowed origins: + for allowed in app_state.ALLOWED_ORIGINS: + try: + if regex.match(origin, allowed): + is_allowed = True + break + except Exception as exception: + app_state.printer(exception) + + # Done here: + app_state.printer(is_allowed) + return is_allowed + + +# --------------------------------------------------------------------------------------------------------------------- + + +async def init( + script_id: str, + debug: bool +): + + """ + To initialize all credentials, instances, and connectivity for this whole script. + :param script_id: The id to use to load cred and data from the internal service. + :param debug: Whether, or not, you would like to print the debug messages. + :return: True if initialized successfully, else False. + """ + + # Declare the required global variables: + # global SCRIPT_DATA + # global ALLOWED_ORIGINS + # global redis_cache + # global kafka_consumer + + # Basic stuff: + if debug: app_state.printer.enable() + app_state.printer("Initializing.") + + # ┏┓ • • + # ┃┃┏┓┓┏┓┓┏┓┏ + # ┗┛┛ ┗┗┫┗┛┗┛ + # ┛ + + response = await app_state.http_client.post( + url = r"https://api.thecaoffice.com/ca/get/title", + headers = {"Origin": "https://thecaoffice.com/"}, + data = { + "domainName": "127.0.0.1:1234", + "screenWidth": 1920, + "screenHeight": 1080 + } + ) + if response.status_code not in [200]: + print("FATAL: ALLOWED ORIGINS NOT FETCHED!") + return False + ALLOWED_ORIGINS = [origin["domain"] for origin in response.json().get("data", {}).get("rs2", [])] + app_state.printer(ALLOWED_ORIGINS) + if len(ALLOWED_ORIGINS) < 1: + print("FATAL: ALLOWED ORIGINS IS EMPTY!") + return False + + # ┏┓ ┓ ┓ ┳┓ + # ┃ ┏┓┏┓┏┫ ┏┓┏┓┏┫ ┃┃┏┓╋┏┓ + # ┗┛┛ ┗ ┗┻ ┗┻┛┗┗┻ ┻┛┗┻┗┗┻ + + # Get the script credentials: + response = await app_state.http_client.get( + url = r"https://nexcom.ditscentre.in/internal/cred/get", + headers = {"X-Script-Id": script_id} + ) + if response.status_code not in [200]: + print("FATAL: SCRIPT CREDENTIALS LOADING FAILED!") + return False + script_cred = response.json().get("data") + + # Get the script data: + response = await app_state.http_client.get( + url = r"https://nexcom.ditscentre.in/internal/data/get", + headers = {"X-Script-Id": script_id} + ) + if response.status_code not in [200]: + print("FATAL: SCRIPT DATA LOADING FAILED!") + return False + app_state.SCRIPT_DATA = response.json().get("data") + + # Done with this step: + app_state.printer("Cred and Data loaded.") + + # ┓┏┓ ┏┓ ┏┓┓• + # ┃┫ ┏┓╋┃┏┏┓ ┃ ┃┓┏┓┏┓╋┏ + # ┛┗┛┗┻┛┛┗┗┻ ┗┛┗┗┗ ┛┗┗┛ + + # Create the consumer that will listen to changes in watchlist: + consumer_creds = script_cred["kafka"]["consumer"] + app_state.kafka_consumer = ConsumerKafka( + topic = consumer_creds["topic"], + bootstrap_servers = consumer_creds["config"]["bootstrapServers"], + security_protocol = consumer_creds["config"].get("securityProtocol", "PLAINTEXT"), + ssl_context = get_ssl_context( + ca_file = consumer_creds["config"].get("caFile"), + cert_file = consumer_creds["config"].get("certFile"), + key_file = consumer_creds["config"].get("keyFile"), + ), + serializer = JSONSerializer(), + debug = debug + ) + if not await app_state.kafka_consumer.connect(): + print("FATAL: KAFKA CONSUMER NOT CREATED!") + return False + app_state.printer("Kafka consumer ready.") + + # ┳┓ ┓• ┏┓ ┓ + # ┣┫┏┓┏┫┓┏ ━━ ┃ ┏┓┏┣┓┏┓ + # ┛┗┗ ┗┻┗┛ ┗┛┗┻┗┛┗┗ + + app_state.redis_cache = AsyncRedisCache( + connection_string = script_cred["redisCache"]["general"]["sentinelJson"], + # connection_string = script_cred["redisCache"]["general"]["connectionString"], + # serializer = JSONSerializer(), + debug = debug, + debug_prefix = "General Cache | " + ) + if not await app_state.redis_cache.connect(): + print("FATAL: REDIS CACHE NOT CREATED!") + return False + app_state.printer("Redis cache ready.") + + # ┳┓ + # ┃┃┏┓┏┓┏┓ + # ┻┛┗┛┛┗┗ + + # If everything went well, we return with success: + app_state.printer("Initialization done.") + return True + + +# --------------------------------------------------------------------------------------------------------------------- + + +def register_namespaces(): + + """ + A quick function that registers all the namespaces to the same Socket.IO app. + """ + + app_state.printer("Registering namespaces.") + sio.register_namespace(TestNamespace("/test")) + app_state.printer("Namespaces registered.") + + +# --------------------------------------------------------------------------------------------------------------------- + + +@sio.on(event = app_state.EVENT_ECHO, namespace = app_state.NAMESPACE_DEFAULT) +async def echo(sid, data): + print(f"Received (default ns) from {sid}: {data}") + await sio.emit("echo", data, to = sid) + + +# ***************************************************************************************************************** +# ***** **** +# *** MAIN PROGRAM *** +# ***** **** +# ***************************************************************************************************************** + + +if __name__ == "__main__": + + app_state.printer("Main.") + + # To get args from the terminal: + import argparse + + # To run the ASGI: + import uvicorn + from multiprocessing import freeze_support + + # Get the config from the command-line: + parser = argparse.ArgumentParser(description = f"SocketIO to serve live market data (and a general passthrough).") + parser.add_argument( + "-w", "--workers", + type = int, + help = "The no. of threads to spin up for this instance!", + default = 2 + ) + parser.add_argument( + "-a", "--host", + type = str, + help = "The host for the app. e.g.: '0.0.0.0' or '127.0.0.1'.", + default = "0.0.0.0" + ) + parser.add_argument( + "-p", "--port", + type = int, + help = "The port no. to bind the app to.", + default = 8080 + ) + parser.add_argument( + "-s", "--script-id", + type = str, + help = "The id of this script (will affect the loaded config)." + ) + parser.add_argument( + "-d", "--debug", + action = "store_true", + help = "Whether, or not, you want to see debugging messages in the terminal.", + default = False + ) + args = parser.parse_args() + + # Note down the config; + os.environ["SCRIPT_ID"] = args.script_id + os.environ["DEBUG"] = str(args.debug) + + # Startup message: + app_state.printer.enable() + app_state.printer(str(args.debug)) + if str(args.debug).lower().find("false") >= 0: app_state.printer.disable() + + # Register all namespaces: + app_state.init_func = init + app_state.origin_check_func = origin_is_allowed + # register_namespaces() + + # Run the gateway: + freeze_support() + uvicorn.run( + app = "app:app", + workers = args.workers, + host = args.host, + port = args.port + ) diff --git a/wsio_v2/test/__init__.py b/wsio_v2/test/__init__.py new file mode 100644 index 0000000..e69de29 diff --git a/wsio_v2/test/echo.py b/wsio_v2/test/echo.py new file mode 100644 index 0000000..db4c95a --- /dev/null +++ b/wsio_v2/test/echo.py @@ -0,0 +1,137 @@ +""" + + AUTHOR: + + Khushal P Soonderji + + DATE: + + Create: Tuesday, 30th Sept., 2025 + + OBJECTIVE: + + A simple namespace to test the Socket.IO app with a simple echo utility. + + REFERENCES: + + N/A + + DOWNLOADS: + + N/A + +""" + + +# ***************************************************************************************************************** +# ***** **** +# *** IMPORT *** +# ***** **** +# ***************************************************************************************************************** + + +# To make sibling directories accessible for imports: +import sys +sys.path.append("../wsio") +sys.path.append("..") + +# System-level activities: +import io +import os +import random + +# My utils: +from utils_v2.string import json +from utils_v2.string import regex +from utils_v2.system import files +from utils_v2.date_time import date_time +from utils_v2.database.async_mongo_v2 import AsyncMongo +from utils_v2.queue.kafka.controllers.async_kafka import ConsumerKafka, get_ssl_context +from utils_v2.cache.async_redis_cache_v3 import AsyncRedisCache +from utils_v2.serialization.json_serializer import JSONSerializer + +# To make HTTP calls: +import httpx + +# To work with date and time: +import datetime +import time + +# Models: +from models.core.user import CoreUserInfoModel +from models.finstitutions.trading.ticks import TradingTick + +# To work with SocketIO: +import socket +import socketio + +# To maintain the app's state: +from wsio_v2.app_state import AppState + +# For asynchronous activities: +import asyncio + +# To work with various datatypes: +from typing import List + +# Debugging: +from icecream import IceCreamDebugger + + +# ***************************************************************************************************************** +# ***** **** +# *** MACROS / ONE-TIME INIT *** +# ***** **** +# ***************************************************************************************************************** + + +# --- Nothing Yet + + +# ***************************************************************************************************************** +# ***** **** +# *** VARIABLES *** +# ***** **** +# ***************************************************************************************************************** + + +# --- Nothing Yet + + +# ***************************************************************************************************************** +# ***** **** +# *** FUNCTIONS *** +# ***** **** +# ***************************************************************************************************************** + + +class TestNamespace(socketio.AsyncNamespace): + + # Internal variables: + app_state: AppState = None + + def __init__(self, namespace): + super().__init__(namespace) + print("Test Namespace Init") + + async def on_connect(self, sid, environ, *args): + print(f"ON CONNECT TEST — SID: {sid}") + + async def on_disconnect(self, sid, reason, *args): + print(f"ON DISCONNECT TEST — SID: {sid}") + + async def on_echo(self, sid, data): + print(f"Received echo from {sid}: {data}") + await self.emit("echo", data, to=sid) + + +# ***************************************************************************************************************** +# ***** **** +# *** MAIN PROGRAM *** +# ***** **** +# ***************************************************************************************************************** + + +if __name__ == "__main__": + + pass