从协议到代码:手把手实现MCP消息通道协议栈
1. 从协议理解到代码落地:为什么MCP值得深究
最近在梳理一些跨进程通信和微服务架构的底层实现,MCP(Message Channel Protocol)这个协议反复被提及。它不像HTTP或gRPC那样广为人知,但在一些特定的、对消息传输的可靠性、顺序性和连接管理有严苛要求的内部系统中,却扮演着核心角色。很多朋友在初次接触MCP时,会觉得协议文档读起来有点抽象,各种状态机、确认机制看着头大,更别提自己动手实现一个简单的客户端或服务端了。这正是我写这篇笔记的原因——光看协议规范不够,必须得把代码敲出来,在调试中观察每一个字节的流动,才能真正理解其设计精髓。
所谓MCP,你可以把它想象成在两个应用之间建立的一条“有保障的专用数据管道”。这条管道不仅负责搬运数据(消息),更重要的是它自己有一套严谨的“交通规则”:如何建立连接(握手)、如何确保每个数据包不丢不重(确认与重传)、如何优雅地处理管道临时中断或永久关闭(连接状态管理)。我们常说的TCP是操作系统内核提供的通用、可靠的字节流管道,而MCP则是在应用层,基于TCP或其他可靠传输层,自己再定义了一套面向“消息”的、更贴近业务逻辑的通信契约。理解并实现它,对于构建健壮的分布式系统、自定义网关或中间件,是项非常扎实的基本功。
在这篇笔记里,我不会重复协议的基础定义,而是直接带大家撸起袖子,通过一个可运行的代码示例,拆解MCP核心的传输过程。我们将实现一个简化但功能完整的MCP协议栈,涵盖连接建立、消息封装与解析、确认机制以及连接保活。你会看到协议字段如何映射到代码结构,状态机如何驱动程序行为,以及在实际编码中会遇到哪些“坑”。无论你是想深入了解某个使用了MCP的现有系统,还是为自己的项目设计通信层,希望这篇结合代码的剖析能给你带来直接、实用的参考。
2. 示例项目的顶层设计:模块划分与职责边界
在开始写具体代码之前,我们先花点时间设计整个示例项目的结构。一个清晰的架构能让我们在实现复杂状态逻辑时保持头脑清醒。我们的目标是实现一个双向通信的MCP协议示例,包含客户端(Initiator)和服务端(Acceptor)角色。为了聚焦于协议本身,我们使用本地回环地址(127.0.0.1)和TCP socket进行模拟。
整个项目将分为以下几个核心模块:
- 协议常量与报文定义模块:这里定义了所有MCP协议中用到的常量,如版本号、消息类型(连接请求、连接确认、数据消息、心跳、断开请求等)、错误码。同时,定义核心的协议报文(Frame)的数据结构,包括头部(Header)和载荷(Payload)。
- 编解码器模块:这是协议栈的“翻译官”。负责将内存中的
Frame对象序列化成可以在网络上传输的字节流(编码),以及将接收到的字节流反序列化回Frame对象(解码)。这里会严格遵循MCP协议的二进制格式规范。 - 连接状态机模块:这是MCP协议逻辑的核心。它定义了连接可能处于的各种状态(如
CLOSED,CONNECTING,ESTABLISHED,DISCONNECTING等),以及在不同状态下,接收到不同类型报文或发生特定事件(如超时)时,应如何转移到下一个状态并执行相应动作(如发送确认报文、通知应用层等)。 - 会话管理模块:它封装了一个完整的MCP连接会话。内部会持有一个状态机实例、一个网络传输通道、发送/接收缓冲区、定时器(用于心跳、超时重传等)。它向上对应用层提供简单的
sendMessage()和onMessage()回调接口,向下管理编解码、状态变迁和网络I/O。 - 客户端与服务端启动模块:这是示例的入口。它们负责创建网络监听器或发起连接,将建立的TCP socket包装进我们实现的MCP会话管理中,并演示基本的消息发送与接收流程。
为什么这么设计?因为MCP协议的本质是一个状态驱动的、基于消息的契约。将状态机独立出来,可以使协议逻辑与网络I/O、业务处理解耦,让每一部分的职责单一,便于测试和维护。编解码器独立,则能确保协议格式的严格性,一旦格式需要调整,影响范围也最小。
注意:为了简化示例,我们暂不考虑窗口流控、消息分片、加密等高级特性,只实现最核心的可靠传输机制。同时,错误处理会相对简单,在实际工业级实现中,需要更完备的异常处理和资源管理。
3. 协议报文的结构化定义与编解码实现
任何协议的第一步,都是定义好“对话的语言”——即报文格式。MCP通常采用二进制协议以保证效率和紧凑性。让我们先定义我们的协议帧(Frame)。
一个基本的MCP帧可以分为固定长度的头部(Header)和可变长度的载荷(Payload)。
头部(Header)设计:我们设计一个12字节的固定头部,包含以下字段:
- Magic Number (2字节):协议魔数,用于快速识别和校验,例如
0x4D43(即 'MC' 的ASCII)。 - Version (1字节):协议版本,主版本号。
- Type (1字节):消息类型。例如:0x01(连接请求 CONNECT),0x02(连接确认 CONNACK),0x03(数据消息 DATA),0x04(确认 ACK),0x05(心跳 PING),0x06(心跳回复 PONG),0x07(断开请求 DISCONNECT)。
- Flags (1字节):标志位,可用于表示压缩、加密、是否为请求/响应等。本例暂简单处理。
- Sequence ID (4字节):序列号,用于标识消息和实现确认机制。对于需要确认的消息(如DATA),接收方需回复带有相同Sequence ID的ACK。
- Payload Length (4字节):载荷数据的长度(单位:字节)。最大长度可根据需要设定。
载荷(Payload)设计:对于不同类型的消息,载荷内容不同:
CONNECT: 可包含客户端标识、协议版本协商信息等。CONNACK: 可包含连接状态(成功/失败)、服务端支持的最大消息长度等。DATA: 就是应用层要传输的实际二进制数据。ACK: 通常为空,或可包含一些元信息。PING/PONG: 通常为空。DISCONNECT: 可包含断开原因码。
在代码中,我们会用一个类来表征这个帧:
# protocol.py import struct from enum import IntEnum from dataclasses import dataclass from typing import Optional class FrameType(IntEnum): CONNECT = 0x01 CONNACK = 0x02 DATA = 0x03 ACK = 0x04 PING = 0x05 PONG = 0x06 DISCONNECT = 0x07 @dataclass class McpHeader: magic: int = 0x4D43 # 'MC' version: int = 1 type: FrameType = FrameType.DATA flags: int = 0 seq_id: int = 0 payload_len: int = 0 # 将头部对象打包成12字节的二进制数据 def pack(self) -> bytes: # struct格式: '>H B B B I I' -> 大端序, 2字节无符号短整型,3个1字节无符号字符,2个4字节无符号整型 return struct.pack('>H B B B I I', self.magic, self.version, self.type, self.flags, self.seq_id, self.payload_len) # 从12字节二进制数据解析出头部对象 @classmethod def unpack(cls, data: bytes) -> 'McpHeader': if len(data) != 12: raise ValueError(f"Header length must be 12 bytes, got {len(data)}") magic, version, ftype, flags, seq_id, payload_len = struct.unpack('>H B B B I I', data) return cls(magic, version, FrameType(ftype), flags, seq_id, payload_len) @dataclass class McpFrame: header: McpHeader payload: Optional[bytes] = None def to_bytes(self) -> bytes: """将整个Frame序列化为字节流,用于发送""" header_bytes = self.header.pack() if self.payload: if len(self.payload) != self.header.payload_len: # 这是一个安全校验,确保头部长度字段与实际载荷一致 raise ValueError(f"Payload length mismatch. Header: {self.header.payload_len}, Actual: {len(self.payload)}") return header_bytes + self.payload else: if self.header.payload_len != 0: raise ValueError(f"Header indicates payload length {self.header.payload_len}, but payload is None") return header_bytes @classmethod def from_bytes(cls, data: bytes) -> 'McpFrame': """从字节流解析出一个完整的Frame。注意:调用者需确保data至少包含完整的头部。""" if len(data) < 12: raise ValueError(f"Data too short for a header, need at least 12 bytes, got {len(data)}") header = McpHeader.unpack(data[:12]) total_frame_len = 12 + header.payload_len if len(data) < total_frame_len: raise ValueError(f"Incomplete frame. Need {total_frame_len} bytes, got {len(data)}") payload = data[12:total_frame_len] if header.payload_len > 0 else None return cls(header, payload)编解码器的关键细节与避坑点:
字节序(Endianness):我们使用了
>(大端序,网络字节序)。这是网络协议的标准做法,确保不同架构的机器能正确解析。如果你在本地测试的两端都是x86小端机,但用了大端序编码,而解码用了小端序,数字会完全错乱。务必在协议文档和代码中明确规定并统一使用网络字节序。长度字段的校验:
Header中的payload_len必须与实际的payload长度严格一致。我们在to_bytes方法中做了校验,这是一个防御性编程的好习惯,能及早发现程序逻辑错误,避免发送错误格式的报文导致对端解析崩溃。帧的完整性判断:
from_bytes方法体现了协议解析的一个经典模式——“偷看”头部。我们先解析出固定的12字节头部,从中得到载荷长度,然后才能判断接收缓冲区中是否已经有一个完整帧的数据。如果不够,说明TCP流中数据还未收全,需要等待更多数据到达。这是处理基于TCP的流式协议时必须实现的“拆包”逻辑。常见的错误是假设一次recv调用就能拿到一个完整帧,这在网络波动或消息较大时几乎必然出错。枚举类型的使用:使用
IntEnum定义消息类型,比直接用数字常量更安全、可读性更好。在解析时通过FrameType(ftype)进行转换,如果收到非法的类型值,会抛出ValueError,便于错误处理。
4. 连接状态机的核心逻辑与实现
MCP协议连接的生命周期由状态机驱动。这是协议行为正确性的保证。我们定义一个简化的状态机,包含以下状态:
- CLOSED: 初始状态或最终状态。连接不存在。
- CONNECTING: 客户端已发送CONNECT请求,等待CONNACK响应。
- ESTABLISHED: 连接已成功建立,可以正常收发数据消息。
- DISCONNECTING: 已发送或收到DISCONNECT请求,正在等待连接彻底关闭。
- ERROR: 发生错误(如协议解析错误、超时等),连接异常终止。
状态迁移由事件触发,主要事件包括:
send_connect,send_data,send_disconnect:应用层主动发起的动作。frame_received:从网络接收到一个完整的MCP帧。timer_expired:某个定时器超时(如连接超时、心跳超时、ACK等待超时)。
下面我们用代码勾勒这个状态机的骨架:
# state_machine.py from enum import Enum from typing import Optional, Callable, Any from protocol import McpFrame, FrameType import time class ConnectionState(Enum): CLOSED = "CLOSED" CONNECTING = "CONNECTING" ESTABLISHED = "ESTABLISHED" DISCONNECTING = "DISCONNECTING" ERROR = "ERROR" class McpStateMachine: def __init__(self, on_state_change: Optional[Callable[[ConnectionState], None]] = None): self.state = ConnectionState.CLOSED self.on_state_change = on_state_change # 用于跟踪已发送但未确认的消息,key为seq_id, value为(发送时间, 重试次数, frame) self.pending_acks = {} self.next_seq_id = 1 # 序列号生成器 # 定时器相关(简化处理,实际需用调度器) self.connect_timer = None self.heartbeat_timer = None def _transition_to(self, new_state: ConnectionState): """状态转移,并触发回调""" if self.state != new_state: old_state = self.state self.state = new_state if self.on_state_change: self.on_state_change(old_state, new_state) def handle_event(self, event: str, frame: Optional[McpFrame] = None, **kwargs) -> Optional[McpFrame]: """ 处理事件,返回一个需要发送的Frame(如果有的话)。 这是状态机的核心分发逻辑。 """ if self.state == ConnectionState.CLOSED: return self._handle_closed(event, frame, **kwargs) elif self.state == ConnectionState.CONNECTING: return self._handle_connecting(event, frame, **kwargs) elif self.state == ConnectionState.ESTABLISHED: return self._handle_established(event, frame, **kwargs) elif self.state == ConnectionState.DISCONNECTING: return self._handle_disconnecting(event, frame, **kwargs) # ERROR状态通常不处理事件,等待重置或销毁 return None def _handle_closed(self, event: str, frame: Optional[McpFrame], **kwargs): if event == 'send_connect': # 应用层请求建立连接 connect_frame = McpFrame( header=McpHeader(type=FrameType.CONNECT, seq_id=self._get_next_seq_id()), payload=b'client_info' # 示例载荷 ) self._transition_to(ConnectionState.CONNECTING) # 启动连接超时定时器(示例,实际需定时器模块) self.connect_timer = time.time() + 10 # 10秒超时 return connect_frame return None def _handle_connecting(self, event: str, frame: Optional[McpFrame], **kwargs): if event == 'frame_received' and frame and frame.header.type == FrameType.CONNACK: # 收到连接确认 # 检查CONNACK的载荷,确认连接是否被接受(简化处理,假设成功) self._transition_to(ConnectionState.ESTABLISHED) # 取消连接超时定时器 self.connect_timer = None # 启动心跳定时器 self._start_heartbeat() # 不需要立即回复帧 return None elif event == 'timer_expired' and kwargs.get('timer_id') == 'connect': # 连接超时 self._transition_to(ConnectionState.ERROR) return None # 在CONNECTING状态下收到其他类型的帧,可能是协议错误 elif event == 'frame_received': self._transition_to(ConnectionState.ERROR) return None def _handle_established(self, event: str, frame: Optional[McpFrame], **kwargs): if event == 'send_data': data_payload = kwargs.get('data') if not data_payload: return None data_frame = McpFrame( header=McpHeader(type=FrameType.DATA, seq_id=self._get_next_seq_id(), payload_len=len(data_payload)), payload=data_payload ) # 将消息加入等待确认队列 self.pending_acks[data_frame.header.seq_id] = (time.time(), 0, data_frame) # 启动该消息的ACK超时定时器(简化) return data_frame elif event == 'frame_received' and frame: if frame.header.type == FrameType.DATA: # 收到数据消息,需要回复ACK ack_frame = McpFrame( header=McpHeader(type=FrameType.ACK, seq_id=frame.header.seq_id) ) # 将数据传递给应用层(通过回调或其他机制,此处略) # self.on_data_received(frame.payload) return ack_frame elif frame.header.type == FrameType.ACK: # 收到ACK,从等待队列中移除对应消息 seq_id = frame.header.seq_id if seq_id in self.pending_acks: del self.pending_acks[seq_id] # 可选:触发应用层消息发送成功的回调 elif frame.header.type == FrameType.PING: # 回复PONG return McpFrame(header=McpHeader(type=FrameType.PONG)) elif frame.header.type == FrameType.DISCONNECT: # 对端请求断开 self._transition_to(ConnectionState.DISCONNECTING) # 可以回复一个DISCONNECT(可选),然后关闭连接 return McpFrame(header=McpHeader(type=FrameType.DISCONNECT)) elif event == 'send_disconnect': # 应用层请求断开 self._transition_to(ConnectionState.DISCONNECTING) return McpFrame(header=McpHeader(type=FrameType.DISCONNECT)) elif event == 'timer_expired' and kwargs.get('timer_id') == 'heartbeat': # 心跳超时,发送PING self._start_heartbeat() # 重置心跳定时器 return McpFrame(header=McpHeader(type=FrameType.PING)) return None def _handle_disconnecting(self, event: str, frame: Optional[McpFrame], **kwargs): # 在DISCONNECTING状态,主要等待底层传输关闭,或进行一些清理工作 # 收到对端的DISCONNECT(如果之前是自己发起的)或超时后,可以转移到CLOSED if event == 'frame_received' and frame and frame.header.type == FrameType.DISCONNECT: self._transition_to(ConnectionState.CLOSED) elif event == 'connection_closed': self._transition_to(ConnectionState.CLOSED) return None def _get_next_seq_id(self) -> int: seq = self.next_seq_id self.next_seq_id += 1 # 处理回绕(简单示例) if self.next_seq_id > 0xFFFFFFFF: self.next_seq_id = 1 return seq def _start_heartbeat(self): # 简化:设置一个下次心跳触发的时间点 self.heartbeat_timer = time.time() + 30 # 30秒后发下一次心跳状态机实现的难点与心得:
状态的纯粹性:状态机只负责管理状态和根据事件决定要发送的协议帧。它不应该直接操作网络socket、启动线程或调用复杂的业务逻辑。它通过返回需要发送的
McpFrame对象,以及通过回调(如on_state_change)来影响外部世界。这种设计使得状态机本身易于单元测试。定时器的整合:在实际系统中,定时器管理是个麻烦事。上面的示例用简单的未来时间点来模拟。在真实实现中,你需要一个集中的定时器调度器,能够添加、取消定时器,并在超时时向状态机发送
timer_expired事件。定时器ID需要与特定任务(如等待某seq_id的ACK、心跳)关联。未确认消息的重传:
pending_acks字典是可靠传输的关键。当状态机在ESTABLISHED状态下收到timer_expired事件,且定时器ID对应某个seq_id时,就需要从pending_acks中取出该消息帧进行重传,并增加重试计数。超过最大重试次数后,应将连接置为ERROR状态。这部分逻辑在上例中省略了,但它是必须实现的。并发与线程安全:如果网络I/O和状态机驱动运行在不同的线程,那么对状态机
handle_event的调用必须是线程安全的,通常需要加锁。一个更清晰的架构是使用单线程的事件循环(如asyncio),将所有事件(网络收到数据、定时器超时、应用层发送请求)都放入一个队列,由事件循环顺序取出并调用状态机处理。
5. 会话管理:粘合协议栈与网络I/O
状态机和编解码器准备好了,我们需要一个“会话”来把它们和真实的网络连接结合起来,并管理整个生命周期。这个会话类(McpSession)是给应用层使用的直接接口。
# session.py import socket import threading import queue import time from typing import Optional, Callable from protocol import McpFrame, McpHeader, FrameType from state_machine import McpStateMachine, ConnectionState class McpSession: def __init__(self, sock: socket.socket, is_initiator: bool): self.sock = sock self.is_initiator = is_initiator self.state_machine = McpStateMachine(on_state_change=self._on_state_change) self.send_queue = queue.Queue() # 用于存放待发送的Frame self.receive_buffer = bytearray() # 接收数据的缓冲区 self.running = False self.recv_thread: Optional[threading.Thread] = None self.send_thread: Optional[threading.Thread] = None self.on_message: Optional[Callable[[bytes], None]] = None # 应用层消息回调 def start(self): """启动会话,开始处理I/O""" self.running = True # 启动接收线程 self.recv_thread = threading.Thread(target=self._recv_loop, daemon=True) self.recv_thread.start() # 启动发送线程 self.send_thread = threading.Thread(target=self._send_loop, daemon=True) self.send_thread.start() # 如果是发起方(客户端),主动发送CONNECT if self.is_initiator: connect_frame = self.state_machine.handle_event('send_connect') if connect_frame: self.send_queue.put(connect_frame) def send_data(self, data: bytes): """应用层调用此方法发送数据""" if self.state_machine.state != ConnectionState.ESTABLISHED: raise ConnectionError("Connection is not established") frame_to_send = self.state_machine.handle_event('send_data', data=data) if frame_to_send: self.send_queue.put(frame_to_send) def disconnect(self): """应用层请求断开连接""" frame_to_send = self.state_machine.handle_event('send_disconnect') if frame_to_send: self.send_queue.put(frame_to_send) # 设置标志,让循环退出 self.running = False def _on_state_change(self, old_state: ConnectionState, new_state: ConnectionState): """状态变化回调""" print(f"[Session] State changed: {old_state} -> {new_state}") if new_state == ConnectionState.CLOSED or new_state == ConnectionState.ERROR: self.running = False self.sock.close() def _recv_loop(self): """接收线程循环:从socket读数据,拼装完整帧,交给状态机处理""" while self.running: try: # 设置超时,以便能响应running标志的变化 self.sock.settimeout(1.0) data = self.sock.recv(4096) if not data: # 对端关闭连接 self.state_machine.handle_event('connection_closed') break self.receive_buffer.extend(data) self._process_buffer() except socket.timeout: continue # 超时是正常的,继续循环检查running标志 except (ConnectionError, OSError) as e: print(f"[RecvLoop] Socket error: {e}") self.state_machine.handle_event('connection_closed') break def _process_buffer(self): """处理接收缓冲区,尝试解析出完整帧""" while len(self.receive_buffer) >= 12: # 至少有一个头部的长度 try: # 尝试解析头部,获取完整帧长度 header = McpHeader.unpack(self.receive_buffer[:12]) total_frame_len = 12 + header.payload_len if len(self.receive_buffer) < total_frame_len: # 缓冲区数据还不够一个完整帧,等待下次接收 break # 提取完整帧数据 frame_data = bytes(self.receive_buffer[:total_frame_len]) # 从缓冲区移除已处理的数据 del self.receive_buffer[:total_frame_len] # 解码成Frame对象 frame = McpFrame.from_bytes(frame_data) # 将接收到的帧交给状态机处理 frame_to_send = self.state_machine.handle_event('frame_received', frame=frame) # 如果状态机要求发送回复帧(如ACK, PONG),放入发送队列 if frame_to_send: self.send_queue.put(frame_to_send) # 如果是数据消息,传递给应用层回调 if frame.header.type == FrameType.DATA and self.on_message: self.on_message(frame.payload) except ValueError as e: # 解析错误,协议格式非法,转移到错误状态 print(f"[ProcessBuffer] Protocol error: {e}") self.state_machine.handle_event('protocol_error') break def _send_loop(self): """发送线程循环:从队列取Frame,编码后通过socket发送""" while self.running: try: frame = self.send_queue.get(timeout=1.0) frame_bytes = frame.to_bytes() # 在实际项目中,这里需要考虑TCP粘包问题吗? # 不需要!因为to_bytes()已经将整个帧(头+载荷)打包成一个完整的字节流。 # TCP是流式协议,保证顺序,但不保证消息边界。我们的“消息边界”就是通过“长度字段”自己定义的。 # 接收方_process_buffer中的“拆包”逻辑正是为了解决这个问题。 # 所以这里直接sendall即可。 self.sock.sendall(frame_bytes) except queue.Empty: continue # 队列为空,继续循环 except (ConnectionError, OSError) as e: print(f"[SendLoop] Socket error: {e}") self.state_machine.handle_event('connection_closed') break会话管理中的核心陷阱与解决方案:
TCP粘包/半包问题:这是网络编程新手最常见的坑。
_process_buffer方法是解决这个问题的经典模式。我们永远不能假设一次recv调用返回的数据就是一个完整的应用层报文(MCP帧)。我们必须维护一个应用层缓冲区,不断累积数据,并尝试从缓冲区头部解析出完整的帧。McpHeader中的payload_len字段就是我们判断“完整性”的关键。永远基于“长度字段”来拆包,而不是依赖特殊分隔符或固定次数recv。双工通信与线程设计:我们使用了独立的接收和发送线程。这是一个简单清晰的模型,但引入了线程同步问题。这里
send_queue是线程安全的queue.Queue,用于协调状态机线程(主线程或事件线程)和发送线程。接收线程将收到的帧交给状态机,状态机可能产生需要发送的帧(如ACK),这个帧也需要放入send_queue。要小心避免死锁和竞态条件。对于更高效的实现,可以考虑使用asyncio单线程异步模型。资源清理:当连接关闭或出错时,必须确保socket被正确关闭,线程被正确终止。示例中通过
running标志和sock.settimeout来让线程优雅退出。在实际代码中,还需要加入join线程等待其结束。_on_state_change回调在进入CLOSED或ERROR状态时关闭socket并停止循环,这是一个集中清理的好地方。超时与重传的管理:示例中简化了定时器。在真实场景中,你需要一个更精细的定时器管理器。例如,对于
pending_acks中的每条消息,都应该有一个独立的超时定时器。当定时器触发时,向状态机发送一个携带seq_id的timer_expired事件。状态机检查该消息是否已被确认,若未确认则重传。重传次数过多则断开连接。心跳定时器同理。
6. 从零构建:客户端与服务端的启动与交互演示
最后,我们编写客户端和服务端的启动脚本,将以上所有模块串联起来,进行一个简单的演示。
服务端(Acceptor)代码:
# server.py import socket import threading from session import McpSession def on_server_message(data: bytes): print(f"[Server] Received application data: {data.decode('utf-8', errors='ignore')}") # 这里可以处理业务逻辑,然后回复客户端(示例中省略) def handle_client_connection(client_sock: socket.socket, client_addr: tuple): print(f"[Server] New connection from {client_addr}") # 服务端是被动接受连接的一方,is_initiator=False session = McpSession(client_sock, is_initiator=False) session.on_message = on_server_message session.start() # 模拟:等待一段时间,然后主动发送一条消息给客户端 import time time.sleep(2) if session.state_machine.state == ConnectionState.ESTABLISHED: session.send_data(b"Hello from Server!") # 保持连接一段时间,观察心跳等 time.sleep(10) session.disconnect() print(f"[Server] Connection with {client_addr} closed.") def start_server(host='127.0.0.1', port=9999): server_sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM) server_sock.setsockopt(socket.SOL_SOCKET, socket.SO_REUSEADDR, 1) server_sock.bind((host, port)) server_sock.listen(5) print(f"[Server] Listening on {host}:{port}") try: while True: client_sock, client_addr = server_sock.accept() # 为每个客户端连接创建一个新线程处理(生产环境建议使用线程池) client_thread = threading.Thread(target=handle_client_connection, args=(client_sock, client_addr)) client_thread.daemon = True client_thread.start() except KeyboardInterrupt: print("\n[Server] Shutting down...") finally: server_sock.close() if __name__ == "__main__": start_server()客户端(Initiator)代码:
# client.py import socket import time from session import McpSession def on_client_message(data: bytes): print(f"[Client] Received application data: {data.decode('utf-8', errors='ignore')}") def start_client(host='127.0.0.1', port=9999): sock = socket.socket(socket.AF_INET, socket.SOCK_STREAM) try: sock.connect((host, port)) print(f"[Client] Connected to {host}:{port}") # 客户端是发起方,is_initiator=True session = McpSession(sock, is_initiator=True) session.on_message = on_client_message session.start() # 等待连接建立 time.sleep(1) if session.state_machine.state != ConnectionState.ESTABLISHED: print("[Client] Failed to establish MCP connection.") return # 连接建立后,发送几条测试消息 for i in range(3): msg = f"Test message {i+1} from client".encode('utf-8') print(f"[Client] Sending: {msg.decode()}") session.send_data(msg) time.sleep(1.5) # 间隔发送,便于观察 # 等待接收服务端的消息,并保持连接以观察心跳 print("[Client] Waiting for server messages and heartbeat...") time.sleep(15) # 主动断开连接 session.disconnect() print("[Client] Disconnected.") except ConnectionRefusedError: print(f"[Client] Could not connect to {host}:{port}") except Exception as e: print(f"[Client] Error: {e}") finally: sock.close() if __name__ == "__main__": start_client()运行与观察:
- 首先在终端运行
python server.py。 - 然后在另一个终端运行
python client.py。 - 观察两个终端的输出。你应该能看到类似以下的过程:
- 客户端打印
Connected to 127.0.0.1:9999。 - 服务端打印
New connection from ('127.0.0.1', xxxxx)。 - 双方打印状态变化
[Session] State changed: CLOSED -> CONNECTING,然后-> ESTABLISHED。 - 客户端发送三条测试消息。
- 服务端收到消息后,打印
[Server] Received application data: Test message 1...,并在2秒后向客户端发送Hello from Server!。 - 客户端收到服务端消息并打印。
- 在连接空闲期间,你会看到状态机定时(示例中约30秒)触发
PING/PONG的日志(需要在代码中添加相应打印)。 - 最后,服务端或客户端断开连接,状态变为
DISCONNECTING->CLOSED。
- 客户端打印
演示中的关键点验证:
- 连接握手:客户端发送
CONNECT,服务端回复CONNACK,状态机迁移到ESTABLISHED。 - 可靠数据传输:客户端发送的
DATA消息带有递增的seq_id。服务端收到后应回复对应的ACK。你可以在_handle_established的ACK处理部分添加日志,以确认这一机制工作正常。 - 心跳保活:在
ESTABLISHED状态空闲一段时间后,状态机应触发timer_expired(心跳)事件,发送PING,对端回复PONG。 - 有序关闭:调用
disconnect()会发送DISCONNECT帧,状态机进入DISCONNECTING,收到对端的DISCONNECT或底层连接关闭事件后,进入CLOSED状态。
通过这个从协议定义到代码实现,再到实际运行的完整示例,MCP协议中那些抽象的概念——连接状态、序列号、确认、心跳——都变成了可以观察、可以调试的具体代码逻辑。在实现过程中,最深的体会是协议设计决定了代码结构。一个定义良好的状态机是复杂协议实现的基石,它能将杂乱的网络事件和业务请求梳理得井井有条。而缓冲区管理和定时器处理则是实现层面最容易出bug的地方,需要格外小心。这个示例虽然简化,但已经勾勒出了一个可靠应用层协议栈的核心骨架,在此基础上增加流量控制、多路复用、加密等特性,思路都是相通的。下次当你再阅读其他协议的RFC文档时,不妨尝试用这种“状态机+编解码+会话管理”的框架去理解它,并思考如何用代码实现,相信会有更深的领悟。
