| 33 | |
| 34 | # 模拟客户端类 |
| 35 | class MSG_Client: |
| 36 | def __init__(self, client_id): |
| 37 | self.client_id = client_id |
| 38 | self.redis_conn = redis_connection() |
| 39 | # self.client_channel = client_id # 接收服务器消息 |
| 40 | self.status_channel = f"{self.client_id}_status" # 状态消息频道 |
| 41 | self.data_channel = self.client_id # 数据消息频道 |
| 42 | self.server_channel = f"server_{client_id}" # 发送给服务器 |
| 43 | self.pubsub_status = self.redis_conn.pubsub() # 状态订阅 |
| 44 | self.pubsub_data = self.redis_conn.pubsub() # 数据订阅 |
| 45 | self.pubsub_status.subscribe(self.status_channel) |
| 46 | self.pubsub_data.subscribe(self.data_channel) |
| 47 | # self.pubsub = self.redis_conn.pubsub() |
| 48 | # self.pubsub.subscribe(self.client_channel) |
| 49 | self.running = True |
| 50 | # 客户端启动时注册自己 |
| 51 | self.registered = False # 注册状态 |
| 52 | self.allocated = False |
| 53 | self.register() |
| 54 | |
| 55 | def register(self): |
| 56 | # 发送注册消息到公共频道 |
| 57 | self.redis_conn.publish(REGISTER_CHANNEL, self.client_id) |
| 58 | print(f"Client {self.client_id} registered itself") |
| 59 | |
| 60 | def listen(self): |
| 61 | print(f"Client {self.client_id} started, listening on {self.status_channel}") |
| 62 | while self.running: |
| 63 | message = self.pubsub_status.get_message(timeout=1.0) |
| 64 | if message and message['type'] == 'message': |
| 65 | data = message['data'].decode('utf-8') |
| 66 | if data == "REGISTERED": |
| 67 | self.registered = True |
| 68 | elif data == "ALLOCATED": |
| 69 | self.allocated = True |
| 70 | print(f"Client {self.client_id} received: {message['data'].decode()}") |
| 71 | time.sleep(0.01) |
| 72 | |
| 73 | def send_message(self, message): |
| 74 | timeout = 5 # 最多等待 5 秒 |
| 75 | start_time = time.time() |
| 76 | # while not self.registered and time.time() - start_time < timeout: |
| 77 | while not self.registered: |
| 78 | time.sleep(1) |
| 79 | if not self.registered: |
| 80 | logging.error(f"Client {self.client_id} failed to register within {timeout} seconds") |
| 81 | return |
| 82 | full_message = f"{message}" |
| 83 | self.redis_conn.publish(self.server_channel, full_message) |
| 84 | print(f"Client {self.client_id} sent: {full_message}") |
| 85 | |
| 86 | def stop(self): |
| 87 | self.running = False |
| 88 | # self.pubsub.unsubscribe(self.client_channel) |
| 89 | self.pubsub_status.unsubscribe(self.status_channel) |
| 90 | self.pubsub_data.unsubscribe(self.data_channel) |
| 91 | # self.redis_conn.close() |
| 92 | print(f"Client {self.client_id} stopped") |