Skip to content

Commit d78dc9c

Browse files
committed
use async buffering for message receipt in chainsync
1 parent bd1b9a4 commit d78dc9c

1 file changed

Lines changed: 123 additions & 44 deletions

File tree

lib/chain_sync.ex

Lines changed: 123 additions & 44 deletions
Original file line numberDiff line numberDiff line change
@@ -29,10 +29,13 @@ defmodule Xander.ChainSync do
2929
rewriting it from future invokations of `c:handle_block/2`
3030
3131
Returning `{:ok, :next_block, new_state}` will request the next block once
32-
it's made available. This is the only valid return value.
32+
it's made available.
33+
34+
Returning `{:ok, :stop}` will close the connection to the node.
3335
"""
3436
@callback handle_rollback(point :: map(), state) ::
3537
{:ok, :next_block, new_state}
38+
| {:ok, :stop}
3639
when state: term(), new_state: term()
3740

3841
alias Xander.ChainSync.Intersection
@@ -59,6 +62,7 @@ defmodule Xander.ChainSync do
5962
:socket,
6063
:network,
6164
queue: :queue.new(),
65+
buffer: <<>>,
6266
state: []
6367
]
6468

@@ -242,8 +246,8 @@ defmodule Xander.ChainSync do
242246
end
243247
end
244248

245-
# Handle socket messages that arrive during the race condition window
246-
# between setting active mode and transitioning to new_blocks state
249+
# Handle socket messages that arrive between setting active mode and transitioning to
250+
# the new_blocks state
247251
def catching_up(
248252
:info,
249253
{_tcp_or_ssl, socket, data},
@@ -254,8 +258,17 @@ defmodule Xander.ChainSync do
254258
)
255259

256260
case handle_socket_data(data, module_state) do
257-
:ok -> {:next_state, :new_blocks, module_state}
258-
result -> result
261+
:ok ->
262+
{:next_state, :new_blocks, module_state}
263+
264+
{:keep_state, new_module_state} ->
265+
{:next_state, :new_blocks, new_module_state}
266+
267+
:keep_state_and_data ->
268+
{:next_state, :new_blocks, module_state}
269+
270+
result ->
271+
result
259272
end
260273
end
261274

@@ -287,55 +300,119 @@ defmodule Xander.ChainSync do
287300
end
288301
end
289302

290-
# Common handler for socket data in both catching_up and new_blocks states
303+
# Common handler for socket data in both catching_up and new_blocks states.
304+
# Uses async buffering to avoid blocking reads on an active socket.
291305
defp handle_socket_data(
292306
data,
293-
%__MODULE__{transport: transport, socket: socket} = module_state
307+
%__MODULE__{buffer: buffer} = module_state
294308
) do
295-
with {:ok, %{payload: payload, size: payload_length}} <- Util.plex(data),
296-
remaining_payload_length = payload_length - byte_size(payload),
297-
{:ok, combined_payload} <-
298-
read_remaining_payload(transport, socket, payload, remaining_payload_length),
299-
{:ok, decoded} <- CSResponse.decode(combined_payload) do
300-
handle_decoded_response(decoded, combined_payload, module_state)
301-
else
309+
# Append incoming data to the buffer
310+
new_buffer = buffer <> data
311+
process_buffer(new_buffer, module_state)
312+
end
313+
314+
# Process buffered data, extracting and handling complete frames.
315+
# Multiplexed frames have an 8-byte header followed by payload.
316+
defp process_buffer(buffer, %__MODULE__{transport: transport, socket: socket} = module_state) do
317+
case extract_frame(buffer) do
318+
{:ok, payload, remaining_buffer} ->
319+
# We have a complete frame, try to decode it
320+
case CSResponse.decode(payload) do
321+
{:ok, decoded} ->
322+
# Process this frame, then continue with remaining buffer
323+
case handle_decoded_response(decoded, module_state) do
324+
{:keep_state, new_module_state} ->
325+
# Update buffer and continue processing any remaining data
326+
new_module_state = %{new_module_state | buffer: remaining_buffer}
327+
328+
if byte_size(remaining_buffer) > 0 do
329+
process_buffer(remaining_buffer, new_module_state)
330+
else
331+
{:keep_state, new_module_state}
332+
end
333+
334+
:ok ->
335+
# AwaitReply case - update buffer in state
336+
if byte_size(remaining_buffer) > 0 do
337+
process_buffer(remaining_buffer, %{module_state | buffer: remaining_buffer})
338+
else
339+
:ok
340+
end
341+
342+
other ->
343+
other
344+
end
345+
346+
{:error, :incomplete_cbor_data} ->
347+
# CBOR data is incomplete but we have a complete frame - this shouldn't happen
348+
# Log and wait for more data
349+
Logger.warning("Incomplete CBOR in complete frame, buffering")
350+
:ok = Transport.setopts(transport, socket, active: :once)
351+
{:keep_state, %{module_state | buffer: buffer}}
352+
353+
{:error, reason} ->
354+
Logger.error("Failed to decode frame: #{inspect(reason)}")
355+
:ok = Transport.setopts(transport, socket, active: :once)
356+
{:keep_state, %{module_state | buffer: <<>>}}
357+
end
358+
359+
:incomplete ->
360+
# Not enough data for a complete frame, buffer and wait for more
361+
:ok = Transport.setopts(transport, socket, active: :once)
362+
{:keep_state, %{module_state | buffer: buffer}}
363+
302364
{:error, reason} ->
303-
Logger.error("Failed to process socket data: #{inspect(reason)}")
365+
Logger.error("Failed to extract frame: #{inspect(reason)}")
304366
:ok = Transport.setopts(transport, socket, active: :once)
305-
:keep_state_and_data
367+
{:keep_state, %{module_state | buffer: <<>>}}
368+
end
369+
end
370+
371+
# Extract a complete multiplexed frame from the buffer.
372+
# Returns {:ok, payload, remaining_buffer} or :incomplete
373+
defp extract_frame(buffer) when byte_size(buffer) < 8, do: :incomplete
374+
375+
defp extract_frame(buffer) do
376+
<<header::binary-size(8), rest::binary>> = buffer
377+
378+
case Util.plex(header) do
379+
{:ok, %{size: payload_length}} ->
380+
if byte_size(rest) >= payload_length do
381+
<<payload::binary-size(payload_length), remaining::binary>> = rest
382+
{:ok, payload, remaining}
383+
else
384+
:incomplete
385+
end
386+
387+
{:error, reason} ->
388+
{:error, reason}
306389
end
307390
end
308391

309-
defp handle_decoded_response(%AwaitReply{}, _payload, %__MODULE__{
392+
defp handle_decoded_response(%AwaitReply{}, %__MODULE__{
310393
transport: transport,
311394
socket: socket
312395
}) do
313396
:ok = Transport.setopts(transport, socket, active: :once)
314397
:ok
315398
end
316399

317-
defp handle_decoded_response(%RollForward{header: header}, _payload, module_state) do
400+
defp handle_decoded_response(%RollForward{header: header}, module_state) do
318401
handle_roll_forward(header, module_state)
319402
end
320403

321-
defp handle_decoded_response(%RollBackward{point: point}, _payload, module_state) do
404+
defp handle_decoded_response(%RollBackward{point: point}, module_state) do
322405
handle_roll_backward(point, module_state)
323406
end
324407

325-
defp handle_decoded_response(
326-
unknown_response,
327-
combined_payload,
328-
%__MODULE__{
329-
transport: transport,
330-
socket: socket,
331-
client_module: client_module,
332-
state: client_state
333-
}
334-
) do
408+
defp handle_decoded_response(unknown_response, %__MODULE__{
409+
transport: transport,
410+
socket: socket
411+
}) do
412+
# Log unknown response and continue waiting for next message asynchronously
335413
Logger.debug("Unknown message: #{inspect(unknown_response)}")
336-
read_next_message_continue(transport, socket, combined_payload, client_module, client_state)
337414
:ok = Transport.setopts(transport, socket, active: :once)
338-
:keep_state_and_data
415+
:ok
339416
end
340417

341418
defp handle_roll_forward(
@@ -436,19 +513,6 @@ defmodule Xander.ChainSync do
436513
end
437514
end
438515

439-
# Helper to read remaining payload bytes from socket, or return current payload if nothing left to read
440-
defp read_remaining_payload(_transport, _socket, current_payload, 0), do: {:ok, current_payload}
441-
442-
defp read_remaining_payload(transport, socket, current_payload, recv_payload_length) do
443-
case Transport.recv(transport, socket, recv_payload_length, @recv_timeout) do
444-
{:ok, additional_payload} ->
445-
{:ok, current_payload <> additional_payload}
446-
447-
{:error, reason} ->
448-
{:error, reason}
449-
end
450-
end
451-
452516
# Helper to read a complete multiplexed message (header + payload) from socket
453517
defp recv_message(transport, socket) do
454518
with {:ok, header_bytes} <- Transport.recv(transport, socket, 8, @recv_timeout),
@@ -488,6 +552,17 @@ defmodule Xander.ChainSync do
488552
:ok
489553
end
490554

555+
{:ok, %AwaitReply{}} ->
556+
Logger.debug("Awaiting reply during continue")
557+
:ok = Transport.setopts(transport, socket, active: :once)
558+
:ok
559+
560+
{:ok, %RollBackward{}} ->
561+
Logger.debug("RollBackward during continue")
562+
# Handle rollback appropriately
563+
:ok = Transport.send(transport, socket, Messages.next_request())
564+
read_until_sync(transport, socket, client_module, state)
565+
491566
{:error, :incomplete_cbor_data} ->
492567
read_next_message_continue(
493568
transport,
@@ -496,6 +571,10 @@ defmodule Xander.ChainSync do
496571
client_module,
497572
state
498573
)
574+
575+
{:error, reason} ->
576+
Logger.error("Decode failed: #{inspect(reason)}")
577+
{:error, reason}
499578
end
500579

501580
{:error, reason} ->

0 commit comments

Comments
 (0)