diff options
| author | Ritabrata Das <[email protected]> | 2025-03-13 11:21:36 +0530 |
|---|---|---|
| committer | Ritabrata Das <[email protected]> | 2025-03-13 11:21:36 +0530 |
| commit | b63fda551d01ce10231b46a657c4f9410a66ef37 (patch) | |
| tree | abcbc1154ad16b2e6d578dde6f6cc580bbaa7337 | |
| parent | b2630ec9a36ec39b69fb72bfcb1f15d7c665e054 (diff) | |
Make proxy faster and handle Denial of service attack
| -rw-r--r-- | proxy.py | 442 | ||||
| -rw-r--r-- | requirements.txt | 3 |
2 files changed, 263 insertions, 182 deletions
@@ -3,7 +3,7 @@ Sakuya AC : Perfect and Elegant Proxy for your YSFlight Server Lisenced under GPLv3 """ -import asyncio +import platform, sys from struct import unpack, pack from lib.parseFlightData import parseFlightData from lib import YSchat, YSviaversion, Player, Aircraft, triggerCommand @@ -50,6 +50,21 @@ logging.basicConfig(level=LOGGING_LEVEL) handler = logging.getLogger().handlers[0] handler.setFormatter(ColoredFormatter("%(levelname)s: %(message)c")) +# we try to setup uvloop for faster io after setting the logging +# its supposed to give faster proxy but we will have to test that + +if platform.system() != "Windows": + try: + import uvloop + uvloop.install() + info("uvloop found! Using uvloop for fast I/O") + except ImportError: + info("uvloop not found but can be installed! Use ``pip install uvloop``, uvloop gives faster I/O on unix systems!") +else: + info("Running on Windows, using standard event loop") + +import asyncio + info("Welcome to Sakuya AC") info("Perfect and Elegant Proxy for your YSFlight Server") info("Lisenced under GPLv3") @@ -60,11 +75,20 @@ plugin_manager = PluginManager(CONNECTED_PLAYERS) # Close Connection async def close_connection(client_writer, server_writer): - # if DISCORD_ENABLED: await discord_send_message(CHANNEL_ID, "has left the server!") - client_writer.close() - server_writer.close() - await client_writer.wait_closed() - await server_writer.wait_closed() + try: + if client_writer and not client_writer.is_closing(): + client_writer.close() + await client_writer.wait_closed() + except Exception as e: + warning(f"Error closing client connection: {e}") + + try: + if server_writer and not server_writer.is_closing(): + server_writer.close() + await server_writer.wait_closed() + except Exception as e: + warning(f"Error closing server connection: {e}") + # Handle client connections async def handle_client(client_reader, client_writer): @@ -72,7 +96,6 @@ async def handle_client(client_reader, client_writer): message_to_server = [] player = Player.Player(message_to_server, message_to_client, client_writer) #Initialise the player. - try: # Connect to the actual server server_reader, server_writer = await asyncio.open_connection(SERVER_HOST, SERVER_PORT) @@ -84,213 +107,268 @@ async def handle_client(client_reader, client_writer): debug("Player object initiated") async def forward(reader, writer, direction, player=player): + try: + while True: + try: + #Test if there are any unsent messages to the client or + # server from other processes. + keep_message = True # Reset this before each loop. + if direction == "client_to_server" and message_to_server and writer and not writer.is_closing(): + writer.write(message_to_server.pop(0)) + await writer.drain() + elif direction == "server_to_client" and message_to_client and writer and not writer.is_closing(): + client_writer.write(message_to_client.pop(0)) + await client_writer.drain() - while True: - try: - #Test if there are any unsent messages to the client or - # server from other processes. - keep_message = True # Reset this before each loop. - if len(message_to_client) > 0: - client_writer.write(message_to_client.pop(0)) - await client_writer.drain() - if len(message_to_server) > 0: - server_writer.write(message_to_server.pop(0)) - await server_writer.drain() - - if not reader.at_eof(): # Connection closed - try: - header = await reader.readexactly(4) # Ensures we always get 4 bytes - except asyncio.IncompleteReadError: - await close_connection(client_writer, server_writer) - break - except ConnectionResetError: - await close_connection(client_writer, server_writer) - break - except Exception as e: - await close_connection(client_writer, server_writer) - critical(f"Error reading header: {e}") - break - - if not header: - break # Connection closed - - length = unpack("I", header)[0] + if not reader.at_eof() or not reader._connection_lost : # Connection closed + try: + header = await reader.readexactly(4) # Ensures we always get 4 bytes + except asyncio.IncompleteReadError: + await close_connection(client_writer, server_writer) + break + except ConnectionResetError: + await close_connection(client_writer, server_writer) + break + except Exception as e: + await close_connection(client_writer, server_writer) + critical(f"Error reading header: {e}") + break + + if not header: + break # Connection closed + + length = unpack("I", header)[0] + + packet = await reader.read(length) + + if not packet: + await close_connection(client_writer, server_writer) + break + + data = header + packet + packet_type = PacketManager().get_packet_type(packet) + if direction == "client_to_server": + debug("C2S" + str(packet_type) + str(player.username)) + debug(data) + + try: + + if packet_type == "FSNETCMD_LOGON": + decode = FSNETCMD_LOGON(packet) + """ + for p in CONNECTED_PLAYERS: + if p.username == decode.username: + client_writer.write(YSchat.message(f"Same username {decode.username} is aldready connected to server! Kicked {ipAddr}")) + data = None + info(f"Same username {decode.username} is aldready connected to server! Kicked {ipAddr}") + CONNECTED_PLAYERS.remove(player) + await close_connection(client_writer, server_writer) + """ + + player.login(decode) + info(f"Player {player.username} connected from {player.ip}") + + if player.version != YSF_VERSION and VIA_VERSION: + info(f"ViaVersion enabled : Porting {player.username} from {player.version} to {YSF_VERSION}") + message_to_client.append(YSchat.message(f"Porting you to YSFlight {YSF_VERSION}, This is currently Experimental")) + message_to_client.append(YSchat.message(f"Please report any bugs to the server admin or join with the correct version")) + data = YSviaversion.genViaVersion(player.username, YSF_VERSION) #TODO: Refactor using the FSNETCMD packet. + writer.write(data) + continue + + elif packet_type == "FSNETCMD_JOINREQUEST": + decode = FSNETCMD_JOINREQUEST(packet) + player.iff = decode.iff + if DISCORD_ENABLED: + asyncio.create_task(discord_send_message(CHANNEL_ID, f"{player.username} has took off in a {decode.aircraft}! 🛫")) - packet = await reader.read(length) + elif packet_type == "FSNETCMD_AIRPLANESTATE": + decode = FSNETCMD_AIRPLANESTATE(packet) + player.aircraft.add_state(decode) - if not packet: - await close_connection(client_writer, server_writer) - break + #TODO: Do we want to convert all this to plugins? Probably not, but there is duplicated functionality + # keep_message = plugin_manager.triggar_hook('on_flight_data', packet, player, message_to_client, message_to_server) + # if not keep_message: + # data = None - data = header + packet - packet_type = PacketManager().get_packet_type(packet) - if direction == "client_to_server": - debug("C2S" + str(packet_type) + str(player.username)) - debug(data) + if player.aircraft.prev_life < player.aircraft.life and player.aircraft.prev_life != -1 and not player.aircraft.just_repaired: + cheatingMsg = YSchat.message(f"{HEALTH_HACK_MESSAGE} by {player.username}") + warning(f"Health hack detected for {player.username}, Connected from {player.ip}") + message_to_server.append(cheatingMsg) - try: + player.aircraft.just_repaired = False - if packet_type == "FSNETCMD_LOGON": - decode = FSNETCMD_LOGON(packet) - """ - for p in CONNECTED_PLAYERS: - if p.username == decode.username: - client_writer.write(YSchat.message(f"Same username {decode.username} is aldready connected to server! Kicked {ipAddr}")) + elif packet_type == "FSNETCMD_UNJOIN": + if DISCORD_ENABLED: + asyncio.create_task(discord_send_message(CHANNEL_ID, f"{player.username} has left the airplane! 🛬")) + player.aircraft.reset() + + elif packet_type == "FSNETCMD_WEAPONCONFIG": + #keep_message = plugin_manager.triggar_hook('on_weapon_config', packet, player, message_to_client, message_to_server) + #if not keep_message: + # data = None + pass + + elif packet_type == "FSNETCMD_TEXTMESSAGE": + msg = FSNETCMD_TEXTMESSAGE(packet) + # keep_message = plugin_manager.triggar_hook('on_chat', packet, player, message_to_client, message_to_server) + # if not keep_message: + # data = None + + finalMsg = (f"{player.username} : {msg.message}") + if msg.message.startswith(PREFIX): data = None - info(f"Same username {decode.username} is aldready connected to server! Kicked {ipAddr}") - CONNECTED_PLAYERS.remove(player) - await close_connection(client_writer, server_writer) - """ - - player.login(decode) - info(f"Player {player.username} connected from {player.ip}") - - if player.version != YSF_VERSION and VIA_VERSION: - info(f"ViaVersion enabled : Porting {player.username} from {player.version} to {YSF_VERSION}") - message_to_client.append(YSchat.message(f"Porting you to YSFlight {YSF_VERSION}, This is currently Experimental")) - message_to_client.append(YSchat.message(f"Please report any bugs to the server admin or join with the correct version")) - data = YSviaversion.genViaVersion(player.username, YSF_VERSION) #TODO: Refactor using the FSNETCMD packet. - writer.write(data) - continue - - elif packet_type == "FSNETCMD_JOINREQUEST": - decode = FSNETCMD_JOINREQUEST(packet) - player.iff = decode.iff - if DISCORD_ENABLED: - asyncio.create_task(discord_send_message(CHANNEL_ID, f"{player.username} has took off in a {decode.aircraft}! 🛫")) - - elif packet_type == "FSNETCMD_AIRPLANESTATE": - decode = FSNETCMD_AIRPLANESTATE(packet) - player.aircraft.add_state(decode) - - #TODO: Do we want to convert all this to plugins? Probably not, but there is duplicated functionality - # keep_message = plugin_manager.triggar_hook('on_flight_data', packet, player, message_to_client, message_to_server) + command = msg.message.split(" ")[0][1:] + asyncio.create_task(triggerCommand.triggerCommand(command, msg.message, player, message_to_client, message_to_server, plugin_manager)) + info(f"Command {command} triggered by {player.username}") + else: + if DISCORD_ENABLED: + # Make it non blocking! + asyncio.create_task(discord_send_message(CHANNEL_ID, finalMsg)) + + elif packet_type == "FSNETCMD_LIST": + # keep_message = plugin_manager.triggar_hook('on_list', packet, player, message_to_client, message_to_server) # if not keep_message: - # data = None + # data = None + pass - if player.aircraft.prev_life < player.aircraft.life and player.aircraft.prev_life != -1 and not player.aircraft.just_repaired: - cheatingMsg = YSchat.message(f"{HEALTH_HACK_MESSAGE} by {player.username}") - warning(f"Health hack detected for {player.username}, Connected from {player.ip}") - message_to_server.append(cheatingMsg) + keep_message = triggerRespectiveHook(packet_type, packet, player, message_to_client, message_to_server, plugin_manager) + if not keep_message: data = None - player.aircraft.just_repaired = False + except Exception as e: + warning(f"Error parsing flight data: {e}", exc_info=True) + traceback.print_exc() # This will display the full traceback - elif packet_type == "FSNETCMD_UNJOIN": - if DISCORD_ENABLED: - asyncio.create_task(discord_send_message(CHANNEL_ID, f"{player.username} has left the airplane! 🛬")) - player.aircraft.reset() - - elif packet_type == "FSNETCMD_WEAPONCONFIG": - #keep_message = plugin_manager.triggar_hook('on_weapon_config', packet, player, message_to_client, message_to_server) - #if not keep_message: - # data = None - pass - - elif packet_type == "FSNETCMD_TEXTMESSAGE": - msg = FSNETCMD_TEXTMESSAGE(packet) - # keep_message = plugin_manager.triggar_hook('on_chat', packet, player, message_to_client, message_to_server) - # if not keep_message: - # data = None - - finalMsg = (f"{player.username} : {msg.message}") - if msg.message.startswith(PREFIX): - data = None - command = msg.message.split(" ")[0][1:] - asyncio.create_task(triggerCommand.triggerCommand(command, msg.message, player, message_to_client, message_to_server, plugin_manager)) - info(f"Command {command} triggered by {player.username}") - else: - if DISCORD_ENABLED: - # Make it non blocking! - asyncio.create_task(discord_send_message(CHANNEL_ID, finalMsg)) + else : + debug("S2C" + str(packet_type) + str(player.username)) + debug(data) - elif packet_type == "FSNETCMD_LIST": - # keep_message = plugin_manager.triggar_hook('on_list', packet, player, message_to_client, message_to_server) - # if not keep_message: - # data = None - pass + #if packet_type == "FSNETCMD_AIRCMD": + # print(FSNETCMD_AIRCMD.decode(packet)) - keep_message = triggerRespectiveHook(packet_type, packet, player, message_to_client, message_to_server, plugin_manager) + keep_message = triggerRespectiveHookServer(packet_type, packet, player, message_to_client, message_to_server, plugin_manager) if not keep_message: data = None - except Exception as e: - warning(f"Error parsing flight data: {e}", exc_info=True) - traceback.print_exc() # This will display the full traceback - - else : - debug("S2C" + str(packet_type) + str(player.username)) - debug(data) - - #if packet_type == "FSNETCMD_AIRCMD": - # print(FSNETCMD_AIRCMD.decode(packet)) - - keep_message = triggerRespectiveHookServer(packet_type, packet, player, message_to_client, message_to_server, plugin_manager) - if not keep_message: data = None - - #Coming from the server to the client - if packet_type == "FSNETCMD_ADDOBJECT": - if player.check_add_object(FSNETCMD_ADDOBJECT(packet)): - info(f"{player.username} has spawned an aircraft") - addSmoke = FSNETCMD_WEAPONCONFIG.addSmoke(player.aircraft.id) - message_to_server.append(addSmoke) - - elif packet_type == "FSNETCMD_AIRCMD": - #Check the configs against the current aircraft - #These come server to client, not the other way around. - command = FSNETCMD_AIRCMD(packet) - player.aircraft.check_command(command) - - elif packet_type == "FSNETCMD_PREPARESIMULATION": - welcomeMsg = YSchat.message(WELCOME_MESSAGE.format(username=player.username)) - message_to_server.append(welcomeMsg) - player.is_a_bot = False - if DISCORD_ENABLED: - asyncio.create_task(discord_send_message(CHANNEL_ID, f"{player.username} has joined the server!")) - - # Forward the packet to the other endpoint if the data packet still exists. - if data: - writer.write(data) - await writer.drain() - except (asyncio.CancelledError, ConnectionResetError, BrokenPipeError) as e: - if not player.connection_closed: - await close_connection(client_writer, server_writer) - if e == BrokenPipeError or ConnectionResetError or asyncio.CancelledError: - info(f"Connection closed by {player.username} : {player.ip}") - else: - warning(f"Connection error during packet forwarding: {e}") - break + #Coming from the server to the client + if packet_type == "FSNETCMD_ADDOBJECT": + if player.check_add_object(FSNETCMD_ADDOBJECT(packet)): + info(f"{player.username} has spawned an aircraft") + addSmoke = FSNETCMD_WEAPONCONFIG.addSmoke(player.aircraft.id) + message_to_server.append(addSmoke) + + elif packet_type == "FSNETCMD_AIRCMD": + #Check the configs against the current aircraft + #These come server to client, not the other way around. + command = FSNETCMD_AIRCMD(packet) + player.aircraft.check_command(command) + + elif packet_type == "FSNETCMD_PREPARESIMULATION": + welcomeMsg = YSchat.message(WELCOME_MESSAGE.format(username=player.username)) + message_to_server.append(welcomeMsg) + player.is_a_bot = False + if DISCORD_ENABLED: + asyncio.create_task(discord_send_message(CHANNEL_ID, f"{player.username} has joined the server!")) + + # Forward the packet to the other endpoint if the data packet still exists. + if data: + writer.write(data) + await writer.drain() + except (asyncio.CancelledError, ConnectionResetError, BrokenPipeError) as e: + if not player.connection_closed: + await close_connection(client_writer, server_writer) + if isinstance(e, (BrokenPipeError, ConnectionResetError, asyncio.CancelledError)): + info(f"Connection closed by {player.username} : {player.ip}") + else: + warning(f"Connection error during packet forwarding: {e}") + break + except Exception as e: + debug(f"Forward function exiting for {player.username} in {direction}") + player.connection_closed = True + await close_connection(client_writer, server_writer) + """ # Start forwarding data between client and server await asyncio.gather( forward(client_reader, server_writer, "client_to_server"), forward(server_reader, client_writer, "server_to_client"), ) + """ + + client_to_server_task = None + server_to_client_task = None + + try: + # Start forwarding data between client and server + client_to_server_task = asyncio.create_task( + forward(client_reader, server_writer, "client_to_server") + ) + server_to_client_task = asyncio.create_task( + forward(server_reader, client_writer, "server_to_client") + ) + + # Wait for either task to complete + done, pending = await asyncio.wait( + [client_to_server_task, server_to_client_task], + return_when=asyncio.FIRST_COMPLETED + ) + + # Cancel pending tasks + for task in pending: + task.cancel() + + # Wait for cancelled tasks to finish + if pending: + await asyncio.wait(pending, return_when=asyncio.ALL_COMPLETED) + + except Exception as e: + warning(f"Error in client handler: {e}") + finally: + # Clean up tasks if they still exist + for task in [t for t in [client_to_server_task, server_to_client_task] if t and not t.done()]: + task.cancel() + except (asyncio.CancelledError, ConnectionResetError, BrokenPipeError) as e: if not isinstance(e, BrokenPipeError): critical(f"Connection error: {e}") finally: - """ try: - #client_writer.close() - #await client_writer.wait_closed() - pass - except Exception as e: - if not isinstance(e, BrokenPipeError): - critical(f"Error closing client connection: {e}") - """ - if not player.connection_closed: - player.connection_closed = True - CONNECTED_PLAYERS.remove(player) - if not player.is_a_bot: - for player in CONNECTED_PLAYERS: - player.streamWriterObject.write(YSchat.message(f"{player.username} has left the server!")) + if player in CONNECTED_PLAYERS: + player.connection_closed = True + CONNECTED_PLAYERS.remove(player) + if not player.is_a_bot: + for p in CONNECTED_PLAYERS: + try: + if p.streamWriterObject and not p.streamWriterObject.is_closing(): + p.streamWriterObject.write(YSchat.message(f"{player.username} has left the server!")) + await p.streamWriterObject.drain() + except Exception as e: + warning(f"Error sending disconnect message: {e}") + if DISCORD_ENABLED: await discord_send_message(CHANNEL_ID, f"{player.username} has left the server!") - await close_connection(client_writer, server_writer) + # Always close connections, regardless of previous state + try: + if client_writer and not client_writer.is_closing(): + client_writer.close() + await client_writer.wait_closed() + except Exception as e: + warning(f"Error closing client connection: {e}") + + try: + if server_writer and not server_writer.is_closing(): + server_writer.close() + await server_writer.wait_closed() + except Exception as e: + warning(f"Error closing server connection: {e}") + + except Exception as e: + critical(f"Error during connection cleanup: {e}") + traceback.print_exc() # Start the proxy server async def start_proxy(): - server = await asyncio.start_server(handle_client, "0.0.0.0", PROXY_PORT) + server = await asyncio.start_server(handle_client, "0.0.0.0", PROXY_PORT, backlog=100) info(f"Proxy server listening on port {PROXY_PORT}") if DISCORD_ENABLED: await asyncio.create_task(monitor_channel(CHANNEL_ID, CONNECTED_PLAYERS)) diff --git a/requirements.txt b/requirements.txt index 859a08f..0d7ea1d 100644 --- a/requirements.txt +++ b/requirements.txt @@ -1,9 +1,12 @@ aiohappyeyeballs==2.4.6 aiohttp==3.11.12 aiosignal==1.3.2 +async-timeout==5.0.1 attrs==25.1.0 frozenlist==1.5.0 idna==3.10 multidict==6.1.0 propcache==0.2.1 +typing_extensions==4.12.2 +uvloop==0.21.0 yarl==1.18.3 |
