Sync from GitHub main #1
@@ -41,29 +41,7 @@ class GnsConnectionAdapter:
|
|||||||
with self._lock:
|
with self._lock:
|
||||||
self._closed = True
|
self._closed = True
|
||||||
|
|
||||||
def sendall(self, framed_payload: bytes) -> None:
|
def _handle_send_result(self, packet_type: str, result: SendResult) -> None:
|
||||||
raw = bytes(framed_payload)
|
|
||||||
with self._lock:
|
|
||||||
if self._closed:
|
|
||||||
raise OSError("GNS connection is closed")
|
|
||||||
if not raw.endswith(b"\n") or raw.count(b"\n") != 1:
|
|
||||||
raise OSError("GNS compatibility adapter received invalid TCP framing")
|
|
||||||
payload = raw[:-1]
|
|
||||||
try:
|
|
||||||
packet = decode_packet(payload)
|
|
||||||
except PacketCodecError as error:
|
|
||||||
raise OSError(f"invalid outbound GNS packet: {error}") from error
|
|
||||||
packet_type = packet["type"]
|
|
||||||
wire_payload = payload
|
|
||||||
if is_snapshot_packet(packet_type):
|
|
||||||
with self._lock:
|
|
||||||
sequence = self._snapshot_counters[packet_type].advance()
|
|
||||||
try:
|
|
||||||
wire_payload = encode_snapshot(packet_type, payload, sequence)
|
|
||||||
except SnapshotEnvelopeError as error:
|
|
||||||
raise ValueError(str(error)) from error
|
|
||||||
encoded = EncodedPacket(packet_type, wire_payload, delivery_for_packet_type(packet_type))
|
|
||||||
result = self.transport.send_encoded(self.connection_id, encoded)
|
|
||||||
if result is SendResult.SENT:
|
if result is SendResult.SENT:
|
||||||
return
|
return
|
||||||
if result is SendResult.DROPPED and is_snapshot_packet(packet_type):
|
if result is SendResult.DROPPED and is_snapshot_packet(packet_type):
|
||||||
@@ -81,6 +59,42 @@ class GnsConnectionAdapter:
|
|||||||
raise ValueError("GNS outbound packet exceeds maximum message size")
|
raise ValueError("GNS outbound packet exceeds maximum message size")
|
||||||
raise OSError("GNS outbound send failed")
|
raise OSError("GNS outbound send failed")
|
||||||
|
|
||||||
|
def send_encoded(self, encoded: EncodedPacket) -> None:
|
||||||
|
with self._lock:
|
||||||
|
if self._closed:
|
||||||
|
raise OSError("GNS connection is closed")
|
||||||
|
|
||||||
|
packet_type = encoded.packet_type
|
||||||
|
expected_delivery = delivery_for_packet_type(packet_type)
|
||||||
|
if encoded.delivery is not expected_delivery:
|
||||||
|
raise OSError("GNS packet delivery metadata does not match protocol policy")
|
||||||
|
|
||||||
|
wire_payload = encoded.payload
|
||||||
|
if is_snapshot_packet(packet_type):
|
||||||
|
with self._lock:
|
||||||
|
sequence = self._snapshot_counters[packet_type].advance()
|
||||||
|
try:
|
||||||
|
wire_payload = encode_snapshot(packet_type, encoded.payload, sequence)
|
||||||
|
except SnapshotEnvelopeError as error:
|
||||||
|
raise ValueError(str(error)) from error
|
||||||
|
|
||||||
|
wire_packet = EncodedPacket(packet_type, wire_payload, encoded.delivery)
|
||||||
|
result = self.transport.send_encoded(self.connection_id, wire_packet)
|
||||||
|
self._handle_send_result(packet_type, result)
|
||||||
|
|
||||||
|
def sendall(self, framed_payload: bytes) -> None:
|
||||||
|
"""Temporary compatibility for callers that still emit TCP line framing."""
|
||||||
|
raw = bytes(framed_payload)
|
||||||
|
if not raw.endswith(b"\n") or raw.count(b"\n") != 1:
|
||||||
|
raise OSError("GNS compatibility adapter received invalid TCP framing")
|
||||||
|
payload = raw[:-1]
|
||||||
|
try:
|
||||||
|
packet = decode_packet(payload)
|
||||||
|
except PacketCodecError as error:
|
||||||
|
raise OSError(f"invalid outbound GNS packet: {error}") from error
|
||||||
|
packet_type = packet["type"]
|
||||||
|
self.send_encoded(EncodedPacket(packet_type, payload, delivery_for_packet_type(packet_type)))
|
||||||
|
|
||||||
def close(self) -> None:
|
def close(self) -> None:
|
||||||
with self._lock:
|
with self._lock:
|
||||||
if self._closed:
|
if self._closed:
|
||||||
|
|||||||
Reference in New Issue
Block a user