From 5a579d49c0bee2ccf3d349b96caaf1569caabbbe Mon Sep 17 00:00:00 2001 From: Nawaz Haider Date: Mon, 12 Jan 2026 21:43:00 +0600 Subject: [PATCH] Refactor main function and related methods to remove async/await, replacing with synchronous calls for improved performance and simplicity --- main.py | 37 +++++++++++++++------------------- test.py | 30 +++++++++++++++++++++++++++ utils/clob_orders.py | 48 +++++++++++++++++++------------------------- utils/orderbook.py | 11 +++++----- utils/tokens.py | 21 +++++++++---------- 5 files changed, 81 insertions(+), 66 deletions(-) create mode 100644 test.py diff --git a/main.py b/main.py index d56ab31..b608364 100644 --- a/main.py +++ b/main.py @@ -1,6 +1,6 @@ import os import gc -import asyncio +import time from utils.logger import setup_logging from utils.tokens import fetch_tokens from utils.orderbook import OrderBook, SIGNALES @@ -22,21 +22,21 @@ from config import ( gc.disable() -async def main(): +def main(): logger = setup_logging() set_cpu_affinity() logger.info("Polymarket HFT Market Maker started") init_global_client() - await asyncio.sleep(2) + time.sleep(2) if not is_client_ready(): logger.error("ClobClient is not ready. Exiting.") return - up_token, down_token, market_slug = await fetch_tokens() + up_token, down_token, market_slug = fetch_tokens() book = OrderBook(up_token, down_token, market_slug) book.start() - await asyncio.sleep(5) # Allow some time for initial order book data + time.sleep(5) # Allow some time for initial order book data market_data = book.get_current_market_data() @@ -56,11 +56,11 @@ async def main(): book.stop() logger.info("Trading session ended. Starting new session.") gc.collect() - await asyncio.sleep(10) + time.sleep(10) reset_trades() - up_token, down_token, market_slug = await fetch_tokens() + up_token, down_token, market_slug = fetch_tokens() book = OrderBook(up_token, down_token, market_slug) - asyncio.create_task(cache_token_trading_infos(book)) + cache_token_trading_infos(book) book.start() market_data = book.get_current_market_data() @@ -84,7 +84,7 @@ async def main(): trading_side = book.last_signal if (trading_side == SIGNALES.UP) and not up_trend: - await place_anchor_and_hedge( + order_ids = place_anchor_and_hedge( up_token, down_token, "UP", @@ -94,12 +94,12 @@ async def main(): ) current_trades = increment_trades() logger.info( - f"Placed UP anchor and hedge orders. Total trades: {current_trades}" + f"Placed UP anchor and hedge orders. Total trades: {current_trades}, Order IDs: {order_ids}" ) - await asyncio.sleep(MIN_DELAY_BETWEEN_TRADES_SECONDS) + time.sleep(MIN_DELAY_BETWEEN_TRADES_SECONDS) elif (trading_side == SIGNALES.DOWN) and up_trend: - await place_anchor_and_hedge( + order_ids = place_anchor_and_hedge( up_token, down_token, "DOWN", @@ -109,21 +109,16 @@ async def main(): ) current_trades = increment_trades() logger.info( - f"Placed DOWN anchor and hedge orders. Total trades: {current_trades}" + f"Placed DOWN anchor and hedge orders. Total trades: {current_trades}, Order IDs: {order_ids}" ) - await asyncio.sleep(MIN_DELAY_BETWEEN_TRADES_SECONDS) + time.sleep(MIN_DELAY_BETWEEN_TRADES_SECONDS) - await asyncio.sleep(0.01) + time.sleep(0.01) if __name__ == "__main__": try: - if os.name == "nt": - asyncio.run(main()) - else: - import uvloop - - uvloop.run(main()) + main() except KeyboardInterrupt: print("\nMarket maker stopped by user") except Exception as e: diff --git a/test.py b/test.py new file mode 100644 index 0000000..5bbdf78 --- /dev/null +++ b/test.py @@ -0,0 +1,30 @@ +import asyncio +from concurrent.futures import ThreadPoolExecutor +import time + +async def test(): + return + + +loop = asyncio.get_event_loop() + +times = [] + +for _ in range(1000): + start = time.time() + with ThreadPoolExecutor(max_workers=2) as executor: + tasks = [ + loop.run_in_executor(executor, test), + loop.run_in_executor( + executor, + test, + ), + ] + await asyncio.gather(*tasks) + + end = time.time() + times.append(end - start) + + + +sum(times) \ No newline at end of file diff --git a/utils/clob_orders.py b/utils/clob_orders.py index c1ef50b..8758aee 100644 --- a/utils/clob_orders.py +++ b/utils/clob_orders.py @@ -12,7 +12,7 @@ from utils.trade_counter import decrement_trades logger = logging.getLogger(__name__) -async def cache_token_trading_infos( +def cache_token_trading_infos( order_book, ) -> None: client = get_client() @@ -26,7 +26,7 @@ async def cache_token_trading_infos( client.get_fee_rate_bps(down_token_id) -async def place_anchor_and_hedge( +def place_anchor_and_hedge( up_token_id, down_token_id, anchor_side, price, size=5, signed_orders_cache=None ): if anchor_side == "UP": @@ -36,31 +36,29 @@ async def place_anchor_and_hedge( anchor_token_id = down_token_id hedge_token_id = up_token_id - loop = asyncio.get_event_loop() with ThreadPoolExecutor(max_workers=2) as executor: - tasks = [ - loop.run_in_executor( - executor, - place_limit_order_sync, - anchor_token_id, - price, - size, - signed_orders_cache, - ), - loop.run_in_executor( - executor, - place_limit_order_sync, - hedge_token_id, - round(1 - price - PROFIT_MARGIN, 2), - size, - signed_orders_cache, - ), - ] - order_ids = await asyncio.gather(*tasks) + future1 = executor.submit( + place_limit_order_sync, + anchor_token_id, + price, + size, + signed_orders_cache, + ) + future2 = executor.submit( + place_limit_order_sync, + hedge_token_id, + round(1 - price - PROFIT_MARGIN, 2), + size, + signed_orders_cache, + ) + + # Wait for both to complete + order_ids = [future1.result(), future2.result()] logger.info( f"Placed anchor and hedge orders: Anchor Token ID={anchor_token_id}, Hedge Token ID={hedge_token_id}, Order IDs={order_ids}" ) + return order_ids def place_limit_order_sync( @@ -93,8 +91,4 @@ def place_limit_order_sync( return None -async def place_limit_order( - token_id: str, price: float, size: int = 5, signed_orders_cache=None -) -> str: - """Async wrapper for backwards compatibility""" - return place_limit_order_sync(token_id, price, size, signed_orders_cache) + diff --git a/utils/orderbook.py b/utils/orderbook.py index f1079fe..7d05ea6 100644 --- a/utils/orderbook.py +++ b/utils/orderbook.py @@ -1,5 +1,4 @@ import os -import asyncio import bisect import time import json @@ -110,7 +109,7 @@ class OrderBook: self.thread.start() self.monitoring_thread = threading.Thread( - target=lambda: asyncio.run(self._continuous_trading_monitor()), daemon=True + target=self._continuous_trading_monitor, daemon=True ) self.monitoring_thread.start() @@ -182,14 +181,14 @@ class OrderBook: "asks": asks, } - async def _continuous_trading_monitor(self): + def _continuous_trading_monitor(self): logger.info("Started continuous trading monitor") while self.monitoring_running: try: market_data = self.get_current_market_data() if not market_data: - await asyncio.sleep(0.1) + time.sleep(0.1) continue micro_vs_mid_bps = market_data["micro_vs_mid_bps"] @@ -205,11 +204,11 @@ class OrderBook: if current_signal and current_signal != self.last_signal: self.last_signal = current_signal - await asyncio.sleep(0.005) + time.sleep(0.005) except Exception as e: logger.error(f"Error in continuous trading monitor: {e}") - await asyncio.sleep(1) + time.sleep(1) logger.info("Stopped continuous trading monitor") diff --git a/utils/tokens.py b/utils/tokens.py index 1043abd..9776621 100644 --- a/utils/tokens.py +++ b/utils/tokens.py @@ -2,7 +2,7 @@ import json import logging from multiprocessing.util import get_logger from typing import Optional, Tuple -import aiohttp +import requests from .slug import get_market_slug from config import GAMMA_API_URL, REQUEST_TIMEOUT @@ -10,7 +10,7 @@ from config import GAMMA_API_URL, REQUEST_TIMEOUT logger = logging.getLogger(__name__) -async def fetch_tokens( +def fetch_tokens( coin: str = "btc", ) -> Tuple[Optional[str], Optional[str], Optional[str]]: @@ -21,17 +21,14 @@ async def fetch_tokens( slug = get_market_slug(coin) url = f"{GAMMA_API_URL}/events/slug/{slug}" - async with aiohttp.ClientSession() as session: - async with session.get( - url, timeout=aiohttp.ClientTimeout(total=REQUEST_TIMEOUT) - ) as response: - if response.status == 200: - data = await response.json() - return _extract_tokens(data, slug) - else: - logger.warning(f"API request failed with status {response.status}") + response = requests.get(url, timeout=REQUEST_TIMEOUT) + if response.status_code == 200: + data = response.json() + return _extract_tokens(data, slug) + else: + logger.warning(f"API request failed with status {response.status_code}") - except aiohttp.ClientError as e: + except requests.exceptions.RequestException as e: logger.error(f"Network error fetching tokens: {e}") except json.JSONDecodeError as e: logger.error(f"Invalid JSON response: {e}")