summaryrefslogtreecommitdiff
diff options
context:
space:
mode:
authorRitabrata Das <[email protected]>2025-03-13 11:21:36 +0530
committerRitabrata Das <[email protected]>2025-03-13 11:21:36 +0530
commitb63fda551d01ce10231b46a657c4f9410a66ef37 (patch)
treeabcbc1154ad16b2e6d578dde6f6cc580bbaa7337
parentb2630ec9a36ec39b69fb72bfcb1f15d7c665e054 (diff)
Make proxy faster and handle Denial of service attack
-rw-r--r--proxy.py442
-rw-r--r--requirements.txt3
2 files changed, 263 insertions, 182 deletions
diff --git a/proxy.py b/proxy.py
index 275dd46..889019b 100644
--- a/proxy.py
+++ b/proxy.py
@@ -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