Skip to content

Commit 5ccf66b

Browse files
authored
feat: Add Transport to ChainSync (#53)
* Refactor ChainSync to use Transport
1 parent 1febc1b commit 5ccf66b

3 files changed

Lines changed: 66 additions & 97 deletions

File tree

CHANGELOG.md

Lines changed: 2 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -22,6 +22,8 @@ transport-specific logic.
2222

2323
- Refactored `Xander.Transaction` to use the new `Xander.Transport` module.
2424

25+
- Refactored `Xander.ChainSync` to use the new `Xander.Transport` module.
26+
2527
### Fixed
2628

2729
- Fix reading multiplexer messages on `Xander.Util.plex/1`. This fix properly

lib/chain_sync.ex

Lines changed: 57 additions & 91 deletions
Original file line numberDiff line numberDiff line change
@@ -1,7 +1,6 @@
11
defmodule Xander.ChainSync do
22
@behaviour :gen_statem
33

4-
@basic_transport_opts [:binary, active: false, send_timeout: 4_000]
54
@active_n2c_versions [9, 10, 11, 12, 13, 14, 15, 16]
65
@recv_timeout 5_000
76

@@ -43,6 +42,7 @@ defmodule Xander.ChainSync do
4342
alias Xander.ChainSync.Response.IntersectFound
4443
alias Xander.ChainSync.Response.RollBackward
4544
alias Xander.ChainSync.Response.RollForward
45+
alias Xander.Transport
4646

4747
# Investigate possibly extracting Handshake into a separate state machine.
4848
alias Xander.Handshake.Proposal
@@ -55,9 +55,7 @@ defmodule Xander.ChainSync do
5555
defstruct [
5656
:client_module,
5757
:sync_from,
58-
:client,
59-
:path,
60-
:port,
58+
:transport,
6159
:socket,
6260
:network,
6361
queue: :queue.new(),
@@ -79,12 +77,12 @@ defmodule Xander.ChainSync do
7977
# optionally given by the client of this module
8078
{sync_from, state} = Keyword.pop(opts, :sync_from, nil)
8179

80+
transport = Transport.new(path: path, port: port, type: type)
81+
8282
data = %__MODULE__{
8383
client_module: client_module,
8484
sync_from: sync_from,
85-
client: transport(type),
86-
path: maybe_local_path(path, type),
87-
port: maybe_local_port(port, type),
85+
transport: transport,
8886
network: network,
8987
socket: nil,
9088
state: state
@@ -116,17 +114,13 @@ defmodule Xander.ChainSync do
116114
def disconnected(
117115
:internal,
118116
:connect,
119-
%__MODULE__{client: client, path: path, port: port} = data
117+
%__MODULE__{transport: transport} = data
120118
) do
121-
Logger.debug("Connecting to #{inspect(path)}")
119+
Logger.debug("Connecting...")
122120

123-
case client.connect(
124-
maybe_parse_path(path),
125-
port,
126-
transport_opts(client, path)
127-
) do
121+
case Transport.connect(transport) do
128122
{:ok, socket} ->
129-
Logger.debug("Connected to #{inspect(path)}")
123+
Logger.debug("Connected")
130124
data = %__MODULE__{data | socket: socket}
131125
actions = [{:next_event, :internal, :establish}]
132126
{:next_state, :connected, data, actions}
@@ -146,7 +140,7 @@ defmodule Xander.ChainSync do
146140
:internal,
147141
:establish,
148142
%__MODULE__{
149-
client: client,
143+
transport: transport,
150144
socket: socket,
151145
network: network
152146
} = data
@@ -155,7 +149,7 @@ defmodule Xander.ChainSync do
155149

156150
version_message = Proposal.version_message(@active_n2c_versions, network)
157151

158-
case propose_handshake(client, socket, version_message) do
152+
case propose_handshake(transport, socket, version_message) do
159153
{:ok, _handshake_response} ->
160154
Logger.debug("Handshake successful")
161155
actions = [{:next_event, :internal, :find_intersection}]
@@ -167,9 +161,9 @@ defmodule Xander.ChainSync do
167161
end
168162
end
169163

170-
defp propose_handshake(client, socket, version_message) do
171-
with :ok <- client.send(socket, version_message),
172-
{:ok, full_response} <- client.recv(socket, 0, @recv_timeout),
164+
defp propose_handshake(transport, socket, version_message) do
165+
with :ok <- Transport.send(transport, socket, version_message),
166+
{:ok, full_response} <- Transport.recv(transport, socket, 0, @recv_timeout),
173167
{:ok, handshake_response} <- HSResponse.validate(full_response) do
174168
{:ok, handshake_response}
175169
else
@@ -186,27 +180,27 @@ defmodule Xander.ChainSync do
186180
# until it reaches the tip of the chain - this happens when the client receives
187181
# a msgAwaitReply response.
188182
@spec catching_up(:internal, :find_intersection | :start_chain_sync, %Xander.ChainSync{
189-
:client => atom(),
183+
:transport => Transport.t(),
190184
:socket => any()
191185
}) ::
192186
:keep_state_and_data
193187
| {:keep_state,
194188
%Xander.ChainSync{
195-
:client => atom(),
189+
:transport => Transport.t(),
196190
:socket => any(),
197191
:sync_from => any()
198192
}, [{any(), any(), any()}, ...]}
199193
| {:next_state, :disconnected | :new_blocks,
200-
%Xander.ChainSync{:client => atom(), :socket => any()}}
194+
%Xander.ChainSync{:transport => Transport.t(), :socket => any()}}
201195
def catching_up(
202196
:internal,
203197
:find_intersection,
204-
%__MODULE__{client: client, socket: socket, sync_from: sync_from} = data
198+
%__MODULE__{transport: transport, socket: socket, sync_from: sync_from} = data
205199
) do
206200
with {:ok, %IntersectionTarget{slot: slot, block_bytes: block_bytes}} <-
207-
Intersection.find_target(client, socket, sync_from),
201+
Intersection.find_target(transport, socket, sync_from),
208202
{:ok, %IntersectFound{}} <-
209-
Intersection.find_intersection(client, socket, slot, block_bytes) do
203+
Intersection.find_intersection(transport, socket, slot, block_bytes) do
210204
# Start the actual chainsync messages
211205
actions = [{:next_event, :internal, :start_chain_sync}]
212206
{:keep_state, data, actions}
@@ -220,32 +214,31 @@ defmodule Xander.ChainSync do
220214
def catching_up(
221215
:internal,
222216
:start_chain_sync,
223-
%__MODULE__{client: client, socket: socket, client_module: client_module} =
217+
%__MODULE__{transport: transport, socket: socket, client_module: client_module} =
224218
data
225219
) do
226220
Logger.debug("starting chainsync")
227221

228-
# TODO: move this stuff into a future Transport module
229-
emit_initial_next_message = fn client, socket ->
230-
with :ok <- client.send(socket, Messages.next_request()),
231-
{:ok, header_bytes} <- client.recv(socket, 8, @recv_timeout),
222+
emit_initial_next_message = fn transport, socket ->
223+
with :ok <- Transport.send(transport, socket, Messages.next_request()),
224+
{:ok, header_bytes} <- Transport.recv(transport, socket, 8, @recv_timeout),
232225
# TODO: use Util.plex!
233226
<<_timestamp::big-32, _mode::1, _protocol_id::15, payload_length::big-16>> <-
234227
header_bytes,
235-
{:ok, payload} <- client.recv(socket, payload_length, @recv_timeout) do
228+
{:ok, payload} <- Transport.recv(transport, socket, payload_length, @recv_timeout) do
236229
CSResponse.decode(payload)
237230
else
238231
{:error, reason} ->
239232
{:error, reason}
240233
end
241234
end
242235

243-
case emit_initial_next_message.(client, socket) do
236+
case emit_initial_next_message.(transport, socket) do
244237
{:ok, %RollBackward{}} ->
245-
:ok = client.send(socket, Messages.next_request())
238+
:ok = Transport.send(transport, socket, Messages.next_request())
246239

247240
# Read the next message
248-
read_until_sync(client, socket, client_module, data)
241+
read_until_sync(transport, socket, client_module, data)
249242
{:next_state, :new_blocks, data}
250243

251244
{:error, reason} ->
@@ -259,7 +252,12 @@ defmodule Xander.ChainSync do
259252
def new_blocks(
260253
:info,
261254
{_tcp_or_ssl, socket, data},
262-
%__MODULE__{client: client, socket: socket, client_module: client_module, state: state} =
255+
%__MODULE__{
256+
transport: transport,
257+
socket: socket,
258+
client_module: client_module,
259+
state: state
260+
} =
263261
module_state
264262
) do
265263
Logger.debug("handling new block")
@@ -274,7 +272,7 @@ defmodule Xander.ChainSync do
274272
{:ok, current_payload}
275273

276274
current_payload, recv_payload_length ->
277-
case client.recv(socket, recv_payload_length, @recv_timeout) do
275+
case Transport.recv(transport, socket, recv_payload_length, @recv_timeout) do
278276
{:ok, additional_payload} ->
279277
{:ok, current_payload <> additional_payload}
280278

@@ -289,7 +287,7 @@ defmodule Xander.ChainSync do
289287

290288
case CSResponse.decode(combined_payload) do
291289
{:ok, %AwaitReply{}} ->
292-
:ok = setopts_lib(client).setopts(socket, active: :once)
290+
:ok = Transport.setopts(transport, socket, active: :once)
293291
:keep_state_and_data
294292

295293
{:ok, %RollForward{header: header}} ->
@@ -303,9 +301,9 @@ defmodule Xander.ChainSync do
303301
state
304302
) do
305303
{:ok, :next_block, new_state} ->
306-
:ok = client.send(socket, Messages.next_request())
304+
:ok = Transport.send(transport, socket, Messages.next_request())
307305

308-
{:ok, data} = client.recv(socket, 8, @recv_timeout)
306+
{:ok, data} = Transport.recv(transport, socket, 8, @recv_timeout)
309307
%{payload: payload, size: payload_length} = Util.plex!(data)
310308
remaining_payload_length = payload_length - byte_size(payload)
311309

@@ -315,7 +313,7 @@ defmodule Xander.ChainSync do
315313
case CSResponse.decode(combined_payload) do
316314
{:ok, %AwaitReply{}} ->
317315
# Response should always be [1] msgAwaitReply
318-
:ok = setopts_lib(client).setopts(socket, active: :once)
316+
:ok = Transport.setopts(transport, socket, active: :once)
319317
{:keep_state, %{module_state | state: new_state}}
320318

321319
error ->
@@ -325,7 +323,7 @@ defmodule Xander.ChainSync do
325323

326324
{:close, new_state} ->
327325
Logger.debug("Disconnecting from node")
328-
:ok = client.close(socket)
326+
:ok = Transport.close(transport, socket)
329327
{:next_state, :disconnected, %{module_state | state: new_state}}
330328
end
331329

@@ -342,8 +340,8 @@ defmodule Xander.ChainSync do
342340
state
343341
) do
344342
{:ok, :next_block, new_state} ->
345-
:ok = client.send(socket, Messages.next_request())
346-
:ok = setopts_lib(client).setopts(socket, active: :once)
343+
:ok = Transport.send(transport, socket, Messages.next_request())
344+
:ok = Transport.setopts(transport, socket, active: :once)
347345
{:keep_state, %{module_state | state: new_state}}
348346

349347
{:ok, :stop} ->
@@ -353,7 +351,7 @@ defmodule Xander.ChainSync do
353351
{:error, _} ->
354352
# If decoding fails, try to read another message
355353
read_next_message_continue(
356-
client,
354+
transport,
357355
socket,
358356
combined_payload,
359357
client_module,
@@ -374,14 +372,14 @@ defmodule Xander.ChainSync do
374372
end
375373

376374
# Helper function to read the next message
377-
defp read_until_sync(client, socket, client_module, state) do
375+
defp read_until_sync(transport, socket, client_module, state) do
378376
# Read the header (8 bytes)
379-
case client.recv(socket, 8, @recv_timeout) do
377+
case Transport.recv(transport, socket, 8, @recv_timeout) do
380378
{:ok, header_bytes} ->
381379
# TODO: use Util.plex!
382380
<<_timestamp::big-32, _mode::1, _protocol_id::15, payload_length::big-16>> = header_bytes
383381

384-
case client.recv(socket, payload_length, @recv_timeout) do
382+
case Transport.recv(transport, socket, payload_length, @recv_timeout) do
385383
{:ok, payload} ->
386384
case CSResponse.decode(payload) do
387385
# When we receive a msgAwaitReply, this means we have reached
@@ -397,7 +395,7 @@ defmodule Xander.ChainSync do
397395
# TODO: address race condition that takes place in case the
398396
# node replies after socket is set to active but before the
399397
# client has transitioned to the new state.
400-
:ok = setopts_lib(client).setopts(socket, active: :once)
398+
:ok = Transport.setopts(transport, socket, active: :once)
401399

402400
{:ok, %RollForward{header: header}} ->
403401
# This is the callback from the client module
@@ -409,18 +407,18 @@ defmodule Xander.ChainSync do
409407
state
410408
) do
411409
{:ok, :next_block, new_state} ->
412-
:ok = client.send(socket, Messages.next_request())
413-
read_until_sync(client, socket, client_module, new_state)
410+
:ok = Transport.send(transport, socket, Messages.next_request())
411+
read_until_sync(transport, socket, client_module, new_state)
414412

415413
{:close, new_state} ->
416414
Logger.debug("Disconnecting from node")
417-
:ok = client.close(socket)
415+
:ok = Transport.close(transport, socket)
418416
{:next_state, :disconnected, new_state}
419417
end
420418

421419
{:error, :incomplete_cbor_data} ->
422420
# If decoding fails, try to read another message
423-
read_next_message_continue(client, socket, payload, client_module, state)
421+
read_next_message_continue(transport, socket, payload, client_module, state)
424422
end
425423

426424
{:error, reason} ->
@@ -435,15 +433,15 @@ defmodule Xander.ChainSync do
435433
end
436434

437435
# Helper function to continue reading if the first attempt fails
438-
defp read_next_message_continue(client, socket, first_payload, client_module, state) do
436+
defp read_next_message_continue(transport, socket, first_payload, client_module, state) do
439437
# Read another header
440-
case client.recv(socket, 8, @recv_timeout) do
438+
case Transport.recv(transport, socket, 8, @recv_timeout) do
441439
{:ok, header_bytes} ->
442440
# TODO: use Util.plex!
443441
<<_timestamp::big-32, _mode::1, _protocol_id::15, payload_length::big-16>> = header_bytes
444442

445443
# Read another payload
446-
case client.recv(socket, payload_length, @recv_timeout) do
444+
case Transport.recv(transport, socket, payload_length, @recv_timeout) do
447445
{:ok, second_payload} ->
448446
# Combine the payloads and try to decode
449447
combined_payload = first_payload <> second_payload
@@ -458,16 +456,16 @@ defmodule Xander.ChainSync do
458456
state
459457
) do
460458
{:ok, :next_block, new_state} ->
461-
:ok = client.send(socket, Messages.next_request())
462-
read_until_sync(client, socket, client_module, new_state)
459+
:ok = Transport.send(transport, socket, Messages.next_request())
460+
read_until_sync(transport, socket, client_module, new_state)
463461

464462
{:ok, :stop} ->
465463
:ok
466464
end
467465

468466
{:error, :incomplete_cbor_data} ->
469467
read_next_message_continue(
470-
client,
468+
transport,
471469
socket,
472470
combined_payload,
473471
client_module,
@@ -486,38 +484,6 @@ defmodule Xander.ChainSync do
486484
end
487485
end
488486

489-
### Helper functions
490-
491-
defp maybe_local_path(path, :socket), do: {:local, path}
492-
defp maybe_local_path(path, _), do: path
493-
494-
defp maybe_local_port(_port, :socket), do: 0
495-
defp maybe_local_port(port, _), do: port
496-
497-
defp maybe_parse_path(path) when is_binary(path) do
498-
uri = URI.parse(path)
499-
~c"#{uri.host}"
500-
end
501-
502-
defp maybe_parse_path(path), do: path
503-
504-
defp transport(:ssl), do: :ssl
505-
defp transport(_), do: :gen_tcp
506-
507-
defp setopts_lib(:ssl), do: :ssl
508-
defp setopts_lib(_), do: :inet
509-
510-
defp transport_opts(:ssl, path),
511-
do:
512-
@basic_transport_opts ++
513-
[
514-
verify: :verify_none,
515-
server_name_indication: ~c"#{path}",
516-
secure_renegotiate: true
517-
]
518-
519-
defp transport_opts(_, _), do: @basic_transport_opts
520-
521487
defmacro __using__(_opts) do
522488
quote do
523489
@behaviour Xander.ChainSync

0 commit comments

Comments
 (0)