跳到主要内容

协议与消息

aiosignalr.protocol

IHubProtocol

JSONProtocolMessagePackProtocol 实现的接口。

  • name —— "json" / "messagepack"
  • version —— 协议版本(2)。
  • transfer_format —— TransferFormat.TEXT / BINARY
  • write_message(message) -> bytes —— 把一条消息序列化为帧。
  • feed(data) -> list[HubMessage | None] —— 从字节流解析完整帧。
  • create_parser() -> IHubProtocol —— 新的连接绑定实例。

JSONProtocol

文本协议。每条消息是一个以 0x1E 结尾的 JSON 对象。

MessagePackProtocol

二进制协议。每条消息是一个前缀为 VarInt 长度的 MessagePack 数组。

握手辅助

  • encode_handshake_request(protocol, version) -> bytes
  • encode_handshake_response(error=None) -> bytes
  • HandshakeParser —— 增量、记录分隔符解析器。
  • parse_handshake(frame)HandshakeRequestMessage | HandshakeResponseMessage

aiosignalr.messages

每种 hub 消息的不可变(frozenslots)数据类:

message_type字段
InvocationMessageINVOCATIONtargetargumentsinvocation_idstream_idsheaders
StreamInvocationMessageSTREAM_INVOCATIONtargetargumentsinvocation_idstream_idsheaders
StreamItemMessageSTREAM_ITEMinvocation_iditem
CompletionMessageCOMPLETIONinvocation_iderrorresulthas_result
CancelInvocationMessageCANCEL_INVOCATIONinvocation_id
PingMessagePING
CloseMessageCLOSEerrorallow_reconnect
AckMessageACKsequence_id
SequenceMessageSEQUENCEsequence_id
HandshakeRequestMessageprotocolversion
HandshakeResponseMessageerror

aiosignalr.message_buffer.MessageBuffer

共享的状态化重连缓冲(见状态化重连)。客户端与服务端 都用它。

MessageBuffer(
protocol: IHubProtocol,
write: Callable[[bytes], Awaitable[None]],
*,
buffer_size: int = 100_000,
ack_interval: float = 1.0,
)
方法说明
async send(message, data)缓冲可跟踪消息并写出
async send_serialized(data)缓冲预编码帧(广播 fanout)
should_process(message) -> bool去重 / 门控入站消息
ack(message)丢弃序号 ≤ ack 的缓冲消息
disconnected()进入重连状态(忽略非 sequence 消息)
async resend()发送 Sequence + 重放缓冲
dispose()取消挂起的 ack 定时器

aiosignalr.framing

  • encode_text_message(data) -> bytes —— 追加 0x1E
  • TextMessageParser(max_size) —— 增量 0x1E 分隔解析器。
  • encode_binary_message(payload) -> bytes —— VarInt 长度前缀。
  • BinaryMessageParser(max_size) —— 增量长度前缀解析器。