现状分析
✅ 链路完全跑通:agent截图 → server收到agent_response并转发,网络业务逻辑正常
问题两点:
- server
wait_closed() 没有捕获Windows平台的ConnectionResetError; - controller断开后,server的
controllers列表清理时机优化; - agent侧:当对端server断开要立刻停止screen_stream任务。
⚠️代码仅限本地学习调试,未经授权捕获屏幕、摄像头属于违法行为。
1、修复版 server.py
import asyncioimport jsonTOKEN = "demo123456"PORT = 8899agents = dict()controllers = []agent_id_counter = 0defjson_dumps(msg): s = json.dumps(msg, ensure_ascii=False)return s + "\n"asyncdefhandle_client(reader: asyncio.StreamReader, writer: asyncio.StreamWriter):global agent_id_counter peername = writer.get_extra_info('peername') client_ip = f"{peername[0]}:{peername[1]}" agent_id = None ctrl_item = None print(f"[SERVER]新客户端接入 {client_ip}")try: buffer = b""whileTrue: data = await reader.read(4096)ifnot data:break buffer += datawhileTrue:ifb"\n"notin buffer:break line_bytes, buffer = buffer.split(b"\n", 1) line = line_bytes.decode("utf-8", errors="ignore").strip()ifnot line:continuetry: msg = json.loads(line)except json.JSONDecodeError as e: print(f"[SERVER]JSON解析失败 {e}, line={line}")continue token = msg.get("token")if token != TOKEN:continue msg_type = msg.get("type")if msg_type == "agent_register": agent_id_counter += 1 agent_id = f"agent_{agent_id_counter}" agents[agent_id] = {"reader": reader, "writer": writer, "ip": client_ip} print(f"[SERVER] Agent上线 {agent_id}{client_ip}") ok_msg = json_dumps({"type": "register_ok", "agent_id": agent_id, "token": TOKEN}) writer.write(ok_msg.encode())await writer.drain() agent_list_msg = json_dumps({"type": "agent_list","agents": [{"aid": k, "ip": v["ip"]} for k, v in agents.items()],"token": TOKEN })for c in controllers:try: c["writer"].write(agent_list_msg.encode())await c["writer"].drain()except Exception:passelif msg_type == "controller_login": ctrl_item = {"reader": reader, "writer": writer} controllers.append(ctrl_item) print(f"[SERVER]监控端登录 {client_ip}") agent_list_msg = json_dumps({"type": "agent_list","agents": [{"aid": k, "ip": v["ip"]} for k, v in agents.items()],"token": TOKEN }) writer.write(agent_list_msg.encode())await writer.drain()elif msg_type == "cmd_to_agent": target_aid = msg.get("target_aid") payload = msg.get("payload") cmd = payload.get("cmd") print(f"[SERVER]收到监控下发命令,目标:{target_aid} cmd:{cmd}")if target_aid in agents: target_writer = agents[target_aid]["writer"] forward_msg = json_dumps({"type": "agent_cmd","token": TOKEN,"payload": payload }) target_writer.write(forward_msg.encode())await target_writer.drain() print(f"[SERVER]命令已转发给 {target_aid}")elif msg_type == "agent_response": print(f"[SERVER]收到agent画面数据 agent_id={msg['agent_id']},转发给全部监控") resp_raw = json_dumps(msg) to_remove = []for idx, c in enumerate(controllers):try: c["writer"].write(resp_raw.encode())await c["writer"].drain()except Exception: to_remove.append(idx)for idx in reversed(to_remove):del controllers[idx]except Exception as e: print(f"[SERVER]handle_client exception:{e}")finally: print(f"[SERVER]客户端断开 {client_ip} agent_id={agent_id}")if agent_id isnotNoneand agent_id in agents:del agents[agent_id] agent_list_msg = json_dumps({"type": "agent_list","agents": [{"aid": k, "ip": v["ip"]} for k, v in agents.items()],"token": TOKEN })for c in controllers:try: c["writer"].write(agent_list_msg.encode())await c["writer"].drain()except Exception:passif ctrl_item isnotNoneand ctrl_item in controllers: controllers.remove(ctrl_item) writer.close()# Windows平台捕获WinError64连接重置异常try:await writer.wait_closed()except (ConnectionResetError, OSError):passasyncdefmain(): server = await asyncio.start_server(handle_client, "0.0.0.0", PORT) print(f"中转服务器启动 0.0.0.0:{PORT} token=demo123456")asyncwith server:await server.serve_forever()if __name__ == "__main__": asyncio.run(main())
2、修复版 agent.py
增加连接断开立刻取消流任务,避免server断开后还疯狂发送画面包
import asyncioimport jsonimport base64import iofrom mss import MSSTOKEN = "demo123456"SERVER_HOST = "127.0.0.1"SERVER_PORT = 8899defjson_dumps(msg):return json.dumps(msg, ensure_ascii=False) + "\n"asyncdefagent_loop():whileTrue:try: reader, writer = await asyncio.open_connection(SERVER_HOST, SERVER_PORT) print("[AGENT]已连接中转服务器,注册...") reg_msg = {"type": "agent_register", "token": TOKEN} writer.write(json_dumps(reg_msg).encode())await writer.drain() aid = None stream_task = Noneasyncdefscreen_capture_to_b64(quality=45):def_sync_capture():with MSS() as sct: monitor = sct.monitors[1] sshot = sct.grab(monitor) buf = io.BytesIO()from PIL import Image img = Image.frombytes("RGB", sshot.size, sshot.bgra, "raw", "BGRX") img.save(buf, format="JPEG", quality=quality, optimize=True)return base64.b64encode(buf.getvalue()).decode("utf-8")returnawait asyncio.to_thread(_sync_capture)asyncdefscreen_stream():nonlocal aidwhileTrue:try: b64 = await screen_capture_to_b64(quality=40) resp = {"type": "agent_response","token": TOKEN,"agent_id": aid,"payload": {"stream_b64": b64} } writer.write(json_dumps(resp).encode())await writer.drain()except Exception as e: print(f"[AGENT‑STREAM] err:{e}")breakawait asyncio.sleep(0.35)asyncdefcamera_capture():nonlocal aidtry:import cv2 cap = cv2.VideoCapture(0)await asyncio.sleep(0.5) ret, frame = cap.read() cap.release()if ret:from PIL import Image img = Image.fromarray(cv2.cvtColor(frame, cv2.COLOR_BGR2RGB)) buf = io.BytesIO() img.save(buf, format="JPEG", quality=50) b64 = base64.b64encode(buf.getvalue()).decode("utf-8") resp = {"type": "agent_response","token": TOKEN,"agent_id": aid,"payload": {"camera_b64": b64} } writer.write(json_dumps(resp).encode())await writer.drain() print("[AGENT]摄像头拍照完成,回传")except Exception as e: print(f"[AGENT‑CAM] error:{e}") buffer = ""whileTrue: data = await reader.read(4096)ifnot data: print("[AGENT]服务端连接断开,停止流任务")if stream_task andnot stream_task.done(): stream_task.cancel()break buffer += data.decode("utf-8", errors="ignore")while"\n"in buffer: line, buffer = buffer.split("\n", 1) line = line.strip()ifnot line:continue msg = json.loads(line) token = msg.get("token")if token != TOKEN:continue msg_type = msg.get("type")if msg_type == "register_ok": aid = msg["agent_id"] print(f"[AGENT]注册成功 agent_id={aid}")elif msg_type == "agent_cmd": payload = msg["payload"] cmd = payload.get("cmd") print(f"[AGENT]收到cmd = {cmd}")if cmd == "screenshot": b64_data = await screen_capture_to_b64(quality=50) resp_msg = {"type": "agent_response","token": TOKEN,"agent_id": aid,"payload": {"screenshot_b64": b64_data} } writer.write(json_dumps(resp_msg).encode())await writer.drain() print("[AGENT]单次截图完成回传")elif cmd == "start_screen_stream":if stream_task isNoneor stream_task.done(): stream_task = asyncio.create_task(screen_stream()) print("[AGENT]启动实时屏幕流")elif cmd == "stop_screen_stream":if stream_task andnot stream_task.done(): stream_task.cancel() stream_task = None print("[AGENT]停止实时屏幕流")elif cmd == "camera_snap": asyncio.create_task(camera_capture()) writer.close()await writer.wait_closed()except Exception as e: print(f"[AGENT]连接异常 {e}, 3s后重连")await asyncio.sleep(3)if __name__ == "__main__": asyncio.run(agent_loop())
3、controller_gui.py(增加窗口关闭时优雅关闭socket)
import asyncioimport jsonimport base64import ioimport tkinter as tkfrom tkinter import ttkfrom PIL import Image, ImageTkTOKEN = "demo123456"SERVER_HOST = "127.0.0.1"SERVER_PORT = 8899defjson_dumps(msg):return json.dumps(msg, ensure_ascii=False) + "\n"classRemoteViewer:def__init__(self, root): self.root = root self.root.title("远程查看器 Controller") self.root.geometry("1100x720") self.root.protocol("WM_DELETE_WINDOW", self.on_close) self.reader = None self.writer = None self.agent_list = [] self.selected_aid = None self._exit_flag = False# UI frm_top = ttk.Frame(root) frm_top.pack(fill=tk.X, padx=6, pady=6) ttk.Label(frm_top, text="Agent列表:").pack(side=tk.LEFT) self.lb_agent = tk.Listbox(frm_top, width=40, height=4) self.lb_agent.pack(side=tk.LEFT, padx=6) self.lb_agent.bind("<<ListboxSelect>>", self.on_select_agent) frm_btn = ttk.Frame(root) frm_btn.pack(fill=tk.X, padx=6, pady=2) ttk.Button(frm_btn, text="单次截图", command=self.btn_screenshot).pack(side=tk.LEFT, padx=4) ttk.Button(frm_btn, text="开始实时屏幕", command=self.btn_start_stream).pack(side=tk.LEFT, padx=4) ttk.Button(frm_btn, text="停止屏幕", command=self.btn_stop_stream).pack(side=tk.LEFT, padx=4) ttk.Button(frm_btn, text="摄像头拍照", command=self.btn_camera).pack(side=tk.LEFT, padx=4) frm_img = ttk.Frame(root) frm_img.pack(fill=tk.BOTH, expand=True, padx=6, pady=6) self.lbl_img = ttk.Label(frm_img, text="等待画面...") self.lbl_img.pack(fill=tk.BOTH, expand=True) self.tk_image = Noneimport threading self.loop = asyncio.new_event_loop() th = threading.Thread(target=self._run_loop, daemon=True) th.start()def_run_loop(self): asyncio.set_event_loop(self.loop) self.loop.run_until_complete(self.network_task())asyncdefnetwork_task(self):try: self.reader, self.writer = await asyncio.open_connection(SERVER_HOST, SERVER_PORT) print("[CTRL]开始连接服务器") login_msg = {"type": "controller_login", "token": TOKEN} self.writer.write(json_dumps(login_msg).encode())await self.writer.drain() print("[CTRL]连接服务器成功,发起登录") buffer = ""whilenot self._exit_flag: data = await self.reader.read(4096)ifnot data:break buffer += data.decode("utf-8", errors="ignore")while"\n"in buffer: line, buffer = buffer.split("\n", 1) line = line.strip()ifnot line:continue msg = json.loads(line) token = msg.get("token")if token != TOKEN:continue print(f"[CTRL]收到消息: {msg}") msg_type = msg.get("type")if msg_type == "agent_list": self.agent_list = msg["agents"] self.root.after(0, self.refresh_agent_list)elif msg_type == "agent_response": payload = msg["payload"] b64 = Noneif"screenshot_b64"in payload: b64 = payload["screenshot_b64"]elif"stream_b64"in payload: b64 = payload["stream_b64"]elif"camera_b64"in payload: b64 = payload["camera_b64"]if b64: self.root.after(0, lambda b=b64: self.show_image(b))except Exception as e: print(f"[CTRL]网络异常 {e}")defrefresh_agent_list(self): self.lb_agent.delete(0, tk.END)for item in self.agent_list: text = f"{item['aid']} | {item['ip']}" self.lb_agent.insert(tk.END, text)defon_select_agent(self, event): idx = self.lb_agent.curselection()ifnot idx:return sel = self.agent_list[idx[0]] self.selected_aid = sel["aid"] print(f"[GUI]选中agent_id = {self.selected_aid}")asyncdefsend_cmd_async(self, cmd): print(f"[send_cmd] 准备发送命令: {cmd}, target_aid={self.selected_aid}")if self.selected_aid isNoneor self.writer isNoneor self._exit_flag:return msg = {"type": "cmd_to_agent","target_aid": self.selected_aid,"payload": {"cmd": cmd},"token": TOKEN } raw = json_dumps(msg).encode() print(f"[CTRL‑DEBUG]发送字节长度={len(raw)}, last_byte={raw[-1:]}") self.writer.write(raw)await self.writer.drain() print(f"[send_cmd]命令 {cmd} 发送完成")defbtn_screenshot(self):if self.selected_aid andnot self._exit_flag: asyncio.run_coroutine_threadsafe(self.send_cmd_async("screenshot"), self.loop)defbtn_start_stream(self):if self.selected_aid andnot self._exit_flag: asyncio.run_coroutine_threadsafe(self.send_cmd_async("start_screen_stream"), self.loop)defbtn_stop_stream(self):if self.selected_aid andnot self._exit_flag: asyncio.run_coroutine_threadsafe(self.send_cmd_async("stop_screen_stream"), self.loop)defbtn_camera(self):if self.selected_aid andnot self._exit_flag: asyncio.run_coroutine_threadsafe(self.send_cmd_async("camera_snap"), self.loop)defshow_image(self, b64_str):try: bin_data = base64.b64decode(b64_str) bio = io.BytesIO(bin_data) im = Image.open(bio) im.thumbnail((1050, 620)) self.tk_image = ImageTk.PhotoImage(im) self.lbl_img.config(image=self.tk_image, text="")except Exception as e: print(f"[GUI]图片渲染失败 {e}")defon_close(self): self._exit_flag = Trueif self.writer: self.writer.close() self.root.destroy()if __name__ == "__main__": win = tk.Tk() app = RemoteViewer(win) win.mainloop()
使用
- 全部保存,关闭旧进程;启动顺序:
server.py → agent.py → controller_gui.py - 复现关闭controller窗口:server不再抛出
WinError64异常;agent感知断开,停止屏幕流发包。
后续优化方向(可选)
- JPEG码率自适应,大分辨率截图base64报文过大;
- server增加最大报文长度限制,防止大图片占满内存;
- 传输改为二进制协议,不再用base64,降低带宽占用。