diff --git a/gamespy_natneg_server.py b/gamespy_natneg_server.py index 2a894d6..35c55c2 100644 --- a/gamespy_natneg_server.py +++ b/gamespy_natneg_server.py @@ -7,6 +7,7 @@ import ctypes import struct import threading import time +import Queue import gamespy.gs_utility as gs_utils import other.utils as utils @@ -45,21 +46,32 @@ class GameSpyNatNegServer(object): # Start natneg server address = ('0.0.0.0', 27901) # accessible to outside connections (use this if you don't know what you're doing) - s = socket.socket(socket.AF_INET, socket.SOCK_DGRAM) - s.bind(address) + self.socket = socket.socket(socket.AF_INET, socket.SOCK_DGRAM) + self.socket.bind(address) + + self.write_queue = Queue.Queue(); logger.log(logging.INFO, "Server is now listening on %s:%s..." % (address[0], address[1])) + threading.Thread(target=self.write_queue_worker).start() while 1: - recv_data, addr = s.recvfrom(2048) + recv_data, addr = self.socket.recvfrom(2048) - packet_thread = threading.Thread(target=self.handle_packet, args=(s, recv_data, addr)) - packet_thread.start() + self.handle_packet(recv_data, addr) - def handle_packet(self, s, recv_data, addr): + def write_queue_send(self, data, address): + time.sleep(0.05) + self.socket.sendto(data, address) + + def write_queue_worker(self): + while 1: + data, address = self.write_queue.get() + threading.Thread(target=self.write_queue_send, args=(data, address)).start() + self.write_queue.task_done() + + def handle_packet(self, recv_data, addr): logger.log(logging.DEBUG, "Connection from %s:%d..." % (addr[0], addr[1])) logger.log(logging.DEBUG, utils.pretty_print_hex(recv_data)) - time.sleep(0.05) # Make sure it's a legal packet if recv_data[0:6] != bytearray([0xfd, 0xfc, 0x1e, 0x66, 0x6a, 0xb2]): @@ -75,7 +87,7 @@ class GameSpyNatNegServer(object): output = bytearray(recv_data[0:14]) output += bytearray([0xff, 0xff, 0x6d, 0x16, 0xb5, 0x7d, 0xea ]) # Checked with Tetris DS, Mario Kart DS, and Metroid Prime Hunters, and this seems to be the standard response to 0x00 output[7] = 0x01 # Initialization response - s.sendto(output, addr) + self.write_queue.put((output, addr)) # Try to connect to the server gameid = utils.get_string(recv_data, 0x15) @@ -126,8 +138,8 @@ class GameSpyNatNegServer(object): output += bytearray([0x42, 0x00]) # Unknown, always seems to be \x42\x00 output[7] = 0x05 - #s.sendto(output, (self.session_list[session_id][client_id]['addr'])) - s.sendto(output, (self.session_list[session_id][client_id]['addr'][0], self.session_list[session_id][client_id]['addr'][1])) + #self.write_queue.put((output, (self.session_list[session_id][client_id]['addr']))) + self.write_queue.put((output, (self.session_list[session_id][client_id]['addr'][0], self.session_list[session_id][client_id]['addr'][1]))) logger.log(logging.DEBUG, "Sent connection request to %s:%d..." % (self.session_list[session_id][client_id]['addr'][0], self.session_list[session_id][client_id]['addr'][1])) logger.log(logging.DEBUG, utils.pretty_print_hex(output)) @@ -154,8 +166,8 @@ class GameSpyNatNegServer(object): output += bytearray([0x42, 0x00]) # Unknown, always seems to be \x42\x00 output[7] = 0x05 - #s.sendto(output, (self.session_list[session_id][client]['addr'])) - s.sendto(output, (self.session_list[session_id][client]['addr'][0], self.session_list[session_id][client]['addr'][1])) + #self.write_queue.put((output, (self.session_list[session_id][client]['addr']))) + self.write_queue.put((output, (self.session_list[session_id][client]['addr'][0], self.session_list[session_id][client]['addr'][1]))) logger.log(logging.DEBUG, "Sent connection request to %s:%d..." % (self.session_list[session_id][client]['addr'][0], self.session_list[session_id][client]['addr'][1])) logger.log(logging.DEBUG, utils.pretty_print_hex(output)) @@ -181,7 +193,7 @@ class GameSpyNatNegServer(object): output += bytearray(recv_data[len(output):]) output[7] = 0x0b - s.sendto(output, addr) + self.write_queue.put((output, addr)) logger.log(logging.DEBUG, "Sent address check response to %s:%d..." % (addr[0], addr[1])) logger.log(logging.DEBUG, utils.pretty_print_hex(output)) @@ -192,7 +204,7 @@ class GameSpyNatNegServer(object): output = bytearray(recv_data) output[7] = 0x02 # ERT Test - s.sendto(output, addr) + self.write_queue.put((output, addr)) logger.log(logging.DEBUG, "Sent natify response to %s:%d..." % (addr[0], addr[1])) logger.log(logging.DEBUG, utils.pretty_print_hex(output)) @@ -204,7 +216,7 @@ class GameSpyNatNegServer(object): output = bytearray(recv_data) output[7] = 0x0e # Report response - s.sendto(recv_data, addr) + self.write_queue.put((recv_data, addr)) else: # Was able to connect logger.log(logging.DEBUG, "Received unknown command %02x from %s:%s..." % (ord(recv_data[7]), addr[0], addr[1])) diff --git a/gamespy_qr_server.py b/gamespy_qr_server.py index d5b4a5a..1e73a3e 100644 --- a/gamespy_qr_server.py +++ b/gamespy_qr_server.py @@ -2,11 +2,13 @@ # Query and Reporting: http://docs.poweredbygamespy.com/wiki/Query_and_Reporting_Overview import logging +import select import socket import struct import threading import time import ctypes +import Queue from multiprocessing.managers import BaseManager @@ -69,6 +71,7 @@ class GameSpyQRServer(object): self.socket = socket.socket(socket.AF_INET, socket.SOCK_DGRAM) self.socket.bind(address) + self.socket.setblocking(0) logger.log(logging.INFO, "Server is now listening on %s:%s..." % (address[0], address[1])) @@ -77,18 +80,30 @@ class GameSpyQRServer(object): server_browser_server_thread = threading.Thread(target=server_browser_server.start) server_browser_server_thread.start() - threading.Timer(1, self.keepalive_check).start() + self.write_queue = Queue.Queue(); + self.db = gs_database.GamespyDatabase() + threading.Thread(target=self.write_queue_worker).start() while 1: - recv_data, address = self.socket.recvfrom(2048) + ready = select.select([self.socket], [], [], 15) - packet_thread = threading.Thread(target=self.handle_packet, args=(self.socket, recv_data, address)) - packet_thread.start() + if ready[0]: + recv_data, address = self.socket.recvfrom(2048) + self.handle_packet(self.socket, recv_data, address) + + self.keepalive_check() + + def write_queue_send(self, data, address): + time.sleep(0.05) + self.socket.sendto(data, address) + + def write_queue_worker(self): + while 1: + data, address = self.write_queue.get() + threading.Thread(target=self.write_queue_send, args=(data, address)).start() + self.write_queue.task_done() def handle_packet(self, socket, recv_data, address): - db = gs_database.GamespyDatabase() - time.sleep(0.05) - # Tetris DS overlay 10 @ 02144184 - Handle responses back to server # Tetris DS overlay 10 @ 02144184 - Handle responses back to server # @@ -199,7 +214,7 @@ class GameSpyQRServer(object): # Send message back to client saying it was accepted packet = bytearray([0xfe, 0xfd, 0x0a]) # Send client registered command packet.extend(session_id_raw) # Get the session ID - socket.sendto(packet, address) + self.write_queue.put((packet, address)) self.log(logging.DEBUG, address, "Sent client registered to %s:%s..." % (address[0], address[1])) else: # Failed the challenge, request another during the next heartbeat @@ -230,7 +245,7 @@ class GameSpyQRServer(object): # The endianness of some server data depends on the endianness of the console, so we must be able # to account for that. self.sessions[session_id].playerid = int(k['dwc_pid']) - profile = db.get_profile_from_profileid(self.sessions[session_id].playerid) + profile = self.db.get_profile_from_profileid(self.sessions[session_id].playerid) if "console" in profile: self.sessions[session_id].console = profile['console'] @@ -248,7 +263,7 @@ class GameSpyQRServer(object): packet.extend(server_challenge) packet.extend('\x00') - socket.sendto(packet, address) + self.write_queue.put((packet, address)) self.log(logging.DEBUG, address, "Sent challenge to %s:%s..." % (address[0], address[1])) self.sessions[session_id].sent_challenge = True @@ -300,7 +315,7 @@ class GameSpyQRServer(object): elif recv_data[0] == '\x09': # Available # Availability check only sent to *.available.gs.nintendowifi.net self.log(logging.DEBUG, address, "Received availability request for '%s' from %s:%s..." % (recv_data[5: -1], address[0], address[1])) - socket.sendto(bytearray([0xfe, 0xfd, 0x09, 0x00, 0x00, 0x00, 0x00]), address) + self.write_queue.put((bytearray([0xfe, 0xfd, 0x09, 0x00, 0x00, 0x00, 0x00]), address)) elif recv_data[0] == '\x0a': # Client Registered # Only sent to client, never received? @@ -311,24 +326,21 @@ class GameSpyQRServer(object): self.log(logging.DEBUG, address, utils.pretty_print_hex(recv_data)) def keepalive_check(self): - while 1: - #self.log(logging.DEBUG, None, "Keep alive check on %d sessions" % (len(self.sessions))) + #self.log(logging.DEBUG, None, "Keep alive check on %d sessions" % (len(self.sessions))) - pruned = [] - for session_id in self.sessions: - now = int(time.time()) - delta = now - self.sessions[session_id].keepalive - timeout = 60 # Remove clients that haven't responded in 60 seconds + pruned = [] + for session_id in self.sessions: + now = int(time.time()) + delta = now - self.sessions[session_id].keepalive + timeout = 60 # Remove clients that haven't responded in 60 seconds - if delta < 0 or delta >= timeout: - pruned.append(session_id) - self.server_manager.delete_server(self.sessions[session_id].gamename, self.sessions[session_id].session) - self.log(logging.DEBUG, None, "Keep alive check removed %s:%s for game %s" % (self.sessions[session_id].address[0], self.sessions[session_id].address[1], self.sessions[session_id].gamename)) + if delta < 0 or delta >= timeout: + pruned.append(session_id) + self.server_manager.delete_server(self.sessions[session_id].gamename, self.sessions[session_id].session) + self.log(logging.DEBUG, None, "Keep alive check removed %s:%s for game %s" % (self.sessions[session_id].address[0], self.sessions[session_id].address[1], self.sessions[session_id].gamename)) - for session_id in pruned: - del self.sessions[session_id] - - time.sleep(15.0) + for session_id in pruned: + del self.sessions[session_id] if __name__ == "__main__": qr_server = GameSpyQRServer()