MCPcopy Create free account
hub / github.com/ModelEngine-Group/flexai / MSG_Client

Class MSG_Client

GPU-Virtual-Service/gpu-remoting/scheduler/msg_queue.py:35–92  ·  view source on GitHub ↗

Source from the content-addressed store, hash-verified

33
34# 模拟客户端类
35class 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")

Callers 1

prepare_gpu_listFunction · 0.85

Calls

no outgoing calls

Tested by

no test coverage detected