Packages
macula
2.1.0
7.1.0
7.0.0
6.0.0
5.2.2
5.2.1
5.2.0
5.1.0
5.0.0
4.8.0
4.7.1
4.7.0
4.6.0
4.5.0
4.4.10
4.4.9
4.4.8
4.4.7
4.4.6
4.4.5
4.4.4
4.4.3
4.4.2
4.4.1
4.4.0
4.3.1
4.3.0
4.2.9
4.2.8
4.2.7
4.2.6
4.2.5
4.2.4
4.2.3
4.2.2
4.2.1
4.2.0
4.1.1
4.1.0
4.0.0
3.16.0
3.15.3
3.15.2
3.15.1
3.14.0
3.13.0
3.12.1
3.12.0
3.11.1
3.11.0
3.10.3
3.10.2
3.10.1
3.9.0
3.8.0
3.7.0
3.5.0
3.4.0
3.3.0
3.2.0
3.1.0
3.0.0
2.1.1
2.1.0
2.0.0
1.5.2
1.5.1
1.4.30
1.4.29
1.4.28
1.4.27
1.4.26
1.4.25
1.4.24
1.4.23
1.4.22
1.4.21
1.4.20
1.4.19
1.4.18
1.4.17
1.4.16
1.4.15
1.4.14
1.4.13
1.4.11
1.4.10
1.4.9
1.4.8
1.4.7
1.4.6
1.4.5
1.4.4
1.4.3
1.4.2
1.4.1
1.4.0
1.3.1
1.3.0
1.2.0
1.1.0
1.0.10
1.0.9
1.0.8
1.0.7
1.0.6
1.0.5
1.0.4
1.0.3
1.0.2
1.0.1
1.0.0
0.48.6
0.48.5
0.48.4
0.48.3
0.48.2
0.48.1
0.48.0
0.47.1
0.47.0
0.46.3
0.46.1
0.46.0
0.45.3
0.45.2
0.45.1
0.45.0
0.44.2
0.44.1
0.44.0
0.43.3
0.43.2
0.43.1
0.43.0
0.42.9
0.42.8
0.42.7
0.42.6
0.42.5
0.42.4
0.42.3
0.42.2
0.42.1
0.42.0
0.41.1
0.41.0
0.40.1
0.40.0
0.39.9
0.39.8
0.39.7
0.39.6
0.39.5
0.39.4
0.39.3
0.39.2
0.39.1
0.39.0
0.38.8
0.38.7
0.38.6
0.38.5
0.38.4
0.38.3
0.38.2
0.38.1
0.38.0
0.37.7
0.37.6
0.37.5
0.37.4
0.37.3
0.37.2
0.37.1
0.37.0
0.36.6
0.36.5
0.36.4
0.36.3
0.36.2
0.36.1
0.36.0
0.35.4
0.35.3
0.35.2
0.35.1
0.35.0
0.34.1
0.34.0
0.33.1
0.33.0
0.32.5
0.32.4
0.32.3
0.32.2
0.32.1
0.32.0
0.31.9
0.31.8
0.31.7
0.31.6
0.31.5
0.31.4
0.31.3
0.31.2
0.31.1
0.31.0
0.30.10
0.30.9
0.30.8
0.30.7
0.30.6
0.30.5
0.30.4
0.30.3
0.30.2
0.30.1
0.30.0
0.29.0
0.28.3
0.28.2
0.28.1
0.28.0
0.27.1
0.27.0
0.26.1
0.26.0
0.25.6
0.25.5
0.25.4
0.25.3
0.25.2
0.25.1
0.25.0
0.24.6
0.24.5
0.24.4
0.24.3
0.24.2
0.24.1
0.24.0
0.23.3
0.23.2
0.23.1
0.23.0
0.22.12
0.22.11
0.22.10
0.22.9
0.22.8
0.22.7
0.22.6
0.22.5
0.22.4
0.22.3
0.22.2
0.22.1
0.22.0
0.21.7
0.21.6
0.21.5
0.21.4
0.21.2
0.21.1
0.21.0
0.20.25
0.20.24
0.20.23
0.20.22
0.20.21
0.20.20
0.20.19
0.20.18
0.20.17
0.20.16
0.20.15
0.20.14
0.20.13
0.20.12
0.20.11
0.20.10
0.20.9
0.20.8
0.20.7
0.20.6
0.20.5
0.20.3
0.20.2
0.20.1
0.20.0
0.19.2
0.19.1
0.19.0
0.18.1
0.18.0
0.17.4
0.17.3
0.17.2
0.17.1
0.17.0
0.16.6
0.16.5
0.16.4
0.16.3
0.16.2
0.16.1
0.16.0
0.15.1
0.15.0
0.14.3
0.14.2
0.14.1
0.14.0
0.12.6
0.12.5
0.12.3
0.11.3
0.10.2
0.10.1
0.10.0
0.9.2
0.9.1
0.9.0
0.8.25
0.8.24
0.8.23
0.8.22
0.8.21
0.8.20
0.8.19
0.8.18
0.8.17
0.8.16
0.8.15
0.8.14
0.8.13
0.8.12
0.8.11
0.8.10
0.8.9
0.8.8
0.8.7
0.8.6
0.8.5
0.8.4
0.8.3
0.8.2
0.8.1
0.8.0
0.7.30
0.7.29
0.7.28
0.7.27
0.7.26
0.7.25
0.7.24
0.7.23
0.7.22
0.7.21
0.7.20
0.7.19
0.7.18
0.7.17
0.7.16
0.7.15
0.7.14
0.7.13
0.7.12
0.7.11
0.7.10
0.7.9
0.7.8
0.7.7
0.7.6
0.7.5
0.7.4
0.7.3
0.7.2
0.7.1
0.7.0
0.6.7
0.6.6
0.6.5
0.6.4
0.6.3
0.6.2
0.6.1
0.6.0
0.5.0
0.4.4
0.4.3
0.4.2
0.4.1
0.4.0
0.3.4
0.3.3
0.3.2
0.3.1
Macula HTTP/3 Mesh SDK — connect, subscribe, publish, call, advertise
Current section
Files
Jump to
Current section
Files
src/macula_dist_system/macula_dist_bridge.erl
%%%-------------------------------------------------------------------
%%% @doc Supervised bridge process for relay distribution tunnels.
%%%
%%% Each tunnel gets one bridge gen_server that owns:
%%% - A BridgeSock (gen_tcp, raw byte pipe)
%%% - A reader process (linked, reads socket -> publishes to relay)
%%% - A writer loop (handle_info, receives from relay -> writes to socket)
%%% - Metrics counters for the tunnel
%%% - A relay subscription (unsubscribed on terminate)
%%%
%%% == Relay Reconnection ==
%%%
%%% The relay_client replays subscriptions on reconnect, so the bridge's
%%% incoming data (tunnel_in messages) resumes automatically. The reader
%%% handles publish failures by retrying — if the relay client is
%%% temporarily disconnected, publishes buffer in its gen_server mailbox
%%% and flush when the QUIC connection is re-established.
%%%
%%% If the relay client PID dies (process crash, not just QUIC drop),
%%% the bridge attempts to re-acquire a new client from persistent_term
%%% and re-subscribe. If no client is available within RECONNECT_TIMEOUT,
%%% the bridge exits and the distribution connection drops.
%%%
%%% Started by `macula_dist_bridge_sup' (simple_one_for_one).
%%% @end
%%%-------------------------------------------------------------------
-module(macula_dist_bridge).
-behaviour(gen_server).
-include_lib("kernel/include/logger.hrl").
-export([start_link/1]).
-export([init/1, handle_info/2, handle_cast/2, handle_call/3, terminate/2]).
-define(BRIDGE_RECV_TIMEOUT, 60000).
-define(BACKPRESSURE_HWM, 64).
-define(RECONNECT_INTERVAL, 2000).
-define(RECONNECT_MAX_ATTEMPTS, 15). %% 15 * 2s = 30s max reconnect window
-define(METRIC_BYTES_OUT, 1).
-define(METRIC_BYTES_IN, 2).
-define(METRIC_MSGS_OUT, 3).
-define(METRIC_MSGS_IN, 4).
-define(PUBLISH_RETRIES, 3).
-record(state, {
bridge_sock :: port(),
tunnel_id :: binary(),
client :: pid(),
client_mon :: reference(),
sub_ref :: reference() | undefined,
key :: binary(),
metrics :: counters:counters_ref(),
reader_pid :: pid() | undefined,
send_topic :: binary(),
recv_topic :: binary(),
reconnect_count :: non_neg_integer()
}).
%%%===================================================================
%%% API
%%%===================================================================
-spec start_link(map()) -> {ok, pid()} | {error, term()}.
start_link(Args) ->
gen_server:start_link(?MODULE, Args, []).
%%%===================================================================
%%% gen_server callbacks
%%%===================================================================
init(#{client := Client, bridge_sock := BridgeSock, send_topic := SendTopic,
recv_topic := RecvTopic, tunnel_id := TunnelId, key := Key,
metrics := Metrics}) ->
process_flag(trap_exit, true),
MonRef = erlang:monitor(process, Client),
{ok, SubRef} = subscribe_to_tunnel(Client, RecvTopic),
%% Don't start reader yet — socket ownership hasn't been transferred.
%% The caller sends {socket_ready} after gen_tcp:controlling_process.
?LOG_INFO("[dist_bridge] Waiting for socket handoff for tunnel ~s", [TunnelId]),
{ok, #state{bridge_sock = BridgeSock, tunnel_id = TunnelId, client = Client,
client_mon = MonRef, sub_ref = SubRef, key = Key, metrics = Metrics,
reader_pid = undefined, send_topic = SendTopic, recv_topic = RecvTopic,
reconnect_count = 0}}.
handle_call(_Request, _From, State) ->
{reply, {error, unknown}, State}.
handle_cast(_Msg, State) ->
{noreply, State}.
%% --- Tunnel data from relay subscription ---
handle_info({tunnel_in, EncData}, #state{key = Key, bridge_sock = BridgeSock,
metrics = Metrics, tunnel_id = TunnelId} = State)
when is_binary(EncData) ->
?LOG_DEBUG("[dist_bridge] tunnel_in ~p bytes for ~s", [byte_size(EncData), TunnelId]),
case decrypt(Key, EncData) of
{ok, Data} ->
counters:add(Metrics, ?METRIC_BYTES_IN, byte_size(Data)),
counters:add(Metrics, ?METRIC_MSGS_IN, 1),
case gen_tcp:send(BridgeSock, Data) of
ok ->
{noreply, State};
{error, Reason} ->
?LOG_WARNING("[dist_bridge] Send error ~p for ~s", [Reason, TunnelId]),
{stop, {send_error, Reason}, State}
end;
{error, decrypt_failed} ->
?LOG_WARNING("[dist_bridge] Decrypt failed for ~s", [TunnelId]),
{noreply, State}
end;
%% --- Socket ownership transferred, start reader ---
handle_info(socket_ready, #state{reader_pid = undefined, client = Client,
bridge_sock = BridgeSock, send_topic = SendTopic,
tunnel_id = TunnelId, key = Key, metrics = Metrics} = State) ->
ReaderPid = start_reader(Client, BridgeSock, SendTopic, TunnelId, Key, Metrics),
?LOG_INFO("[dist_bridge] Reader started (~p) for ~s", [ReaderPid, TunnelId]),
{noreply, State#state{reader_pid = ReaderPid}};
%% --- Reader exited (linked process) ---
%% Clean exits (peer closed the dist tunnel, `global` disconnected, reader
%% finished draining, supervisor shutdown) are not crashes — `{stop, normal}`
%% so the gen_server terminates without a CRASH REPORT.
handle_info({'EXIT', Pid, normal}, #state{reader_pid = Pid,
tunnel_id = TunnelId} = State) ->
?LOG_INFO("[dist_bridge] Reader closed cleanly for ~s", [TunnelId]),
{stop, normal, State};
handle_info({'EXIT', Pid, shutdown}, #state{reader_pid = Pid,
tunnel_id = TunnelId} = State) ->
?LOG_INFO("[dist_bridge] Reader shutdown for ~s", [TunnelId]),
{stop, normal, State};
handle_info({'EXIT', Pid, {shutdown, _}}, #state{reader_pid = Pid,
tunnel_id = TunnelId} = State) ->
?LOG_INFO("[dist_bridge] Reader shutdown for ~s", [TunnelId]),
{stop, normal, State};
handle_info({'EXIT', Pid, Reason}, #state{reader_pid = Pid,
tunnel_id = TunnelId} = State) ->
?LOG_WARNING("[dist_bridge] Reader crashed (~p) for ~s", [Reason, TunnelId]),
{stop, {reader_exit, Reason}, State};
handle_info({'EXIT', _Pid, _Reason}, State) ->
{noreply, State};
%% --- Relay client died (Gap 2: reconnection) ---
handle_info({'DOWN', MonRef, process, _Pid, Reason},
#state{client_mon = MonRef, tunnel_id = TunnelId} = State) ->
?LOG_WARNING("[dist_bridge] Relay client died (~p) for ~s, attempting reconnect",
[Reason, TunnelId]),
attempt_reconnect(State#state{client_mon = undefined, sub_ref = undefined,
reconnect_count = 0});
%% --- Reconnect timer ---
handle_info(reconnect_tick, State) ->
attempt_reconnect(State);
handle_info(_Msg, State) ->
{noreply, State}.
terminate(_Reason, #state{bridge_sock = BridgeSock, tunnel_id = TunnelId,
client = Client, sub_ref = SubRef,
client_mon = MonRef}) ->
?LOG_INFO("[dist_bridge] Cleaning up tunnel ~s", [TunnelId]),
catch demonitor_if_set(MonRef),
catch unsubscribe_if_set(Client, SubRef),
catch gen_tcp:close(BridgeSock),
remove_metrics(TunnelId),
ok.
%%%===================================================================
%%% Internal — Reconnection (Gap 2)
%%%===================================================================
attempt_reconnect(#state{reconnect_count = N, tunnel_id = TunnelId} = State)
when N >= ?RECONNECT_MAX_ATTEMPTS ->
?LOG_ERROR("[dist_bridge] Reconnect exhausted (~p attempts) for ~s", [N, TunnelId]),
{stop, relay_reconnect_failed, State};
attempt_reconnect(#state{recv_topic = RecvTopic, tunnel_id = TunnelId,
reconnect_count = N} = State) ->
case macula_dist_relay:get_mesh_client() of
undefined ->
?LOG_INFO("[dist_bridge] No relay client yet, retry ~p for ~s", [N + 1, TunnelId]),
erlang:send_after(?RECONNECT_INTERVAL, self(), reconnect_tick),
{noreply, State#state{reconnect_count = N + 1}};
NewClient ->
?LOG_INFO("[dist_bridge] Re-acquired relay client for ~s", [TunnelId]),
MonRef = erlang:monitor(process, NewClient),
{ok, SubRef} = subscribe_to_tunnel(NewClient, RecvTopic),
%% Restart reader with new client
catch exit(State#state.reader_pid, kill),
NewReader = start_reader(NewClient, State#state.bridge_sock,
State#state.send_topic, TunnelId,
State#state.key, State#state.metrics),
{noreply, State#state{client = NewClient, client_mon = MonRef,
sub_ref = SubRef, reader_pid = NewReader,
reconnect_count = 0}}
end.
%%%===================================================================
%%% Internal — Reader (linked process with publish retry)
%%%===================================================================
start_reader(Client, BridgeSock, SendTopic, TunnelId, Key, Metrics) ->
spawn_link(fun() ->
bridge_reader_loop(Client, BridgeSock, SendTopic, TunnelId, Key, Metrics)
end).
%% Read from the loopback socket and forward to the relay.
%% {error, timeout} is NOT fatal — just means no traffic in the window.
%% Only exit on {error, closed} or genuine errors.
bridge_reader_loop(MeshClient, BridgeSock, SendTopic, TunnelId, Key, Metrics) ->
case gen_tcp:recv(BridgeSock, 0, ?BRIDGE_RECV_TIMEOUT) of
{ok, Data} ->
maybe_backpressure(MeshClient),
Encrypted = encrypt(Key, Data),
publish_with_retry(MeshClient, SendTopic, Encrypted, ?PUBLISH_RETRIES),
counters:add(Metrics, ?METRIC_BYTES_OUT, byte_size(Data)),
counters:add(Metrics, ?METRIC_MSGS_OUT, 1),
bridge_reader_loop(MeshClient, BridgeSock, SendTopic, TunnelId, Key, Metrics);
{error, timeout} ->
%% No traffic in the window — keep reading.
bridge_reader_loop(MeshClient, BridgeSock, SendTopic, TunnelId, Key, Metrics);
{error, closed} ->
?LOG_INFO("[dist_bridge] Reader closed for ~s", [TunnelId]);
{error, Reason} ->
?LOG_WARNING("[dist_bridge] Reader error ~p for ~s", [Reason, TunnelId])
end.
publish_with_retry(_Client, _Topic, _Data, 0) ->
?LOG_WARNING("[dist_bridge] Publish retries exhausted"),
ok;
publish_with_retry(Client, Topic, Data, Retries) ->
case is_process_alive(Client) of
true ->
macula_mesh_client:publish(Client, Topic, Data);
false ->
timer:sleep(1000),
publish_with_retry(Client, Topic, Data, Retries - 1)
end.
subscribe_to_tunnel(Client, RecvTopic) ->
Self = self(),
macula_mesh_client:subscribe(Client, RecvTopic,
fun(Msg) -> Self ! {tunnel_in, macula_dist_relay:extract_payload(Msg)} end).
%%%===================================================================
%%% Internal — Crypto
%%%===================================================================
encrypt(Key, Plaintext) ->
Nonce = crypto:strong_rand_bytes(12),
{Ciphertext, Tag} = crypto:crypto_one_time_aead(
aes_256_gcm, Key, Nonce, Plaintext, <<>>, true),
<<Nonce/binary, Tag/binary, Ciphertext/binary>>.
decrypt(Key, <<Nonce:12/binary, Tag:16/binary, Ciphertext/binary>>) ->
case crypto:crypto_one_time_aead(
aes_256_gcm, Key, Nonce, Ciphertext, <<>>, Tag, false) of
error -> {error, decrypt_failed};
Plaintext -> {ok, Plaintext}
end;
decrypt(_Key, _Data) ->
{error, decrypt_failed}.
%%%===================================================================
%%% Internal — Helpers
%%%===================================================================
maybe_backpressure(MeshClient) ->
case erlang:process_info(MeshClient, message_queue_len) of
{message_queue_len, Len} when Len > ?BACKPRESSURE_HWM ->
timer:sleep(1);
_ ->
ok
end.
demonitor_if_set(undefined) -> ok;
demonitor_if_set(Ref) -> erlang:demonitor(Ref, [flush]).
unsubscribe_if_set(_Client, undefined) -> ok;
unsubscribe_if_set(Client, Ref) -> macula_mesh_client:unsubscribe(Client, Ref).
remove_metrics(TunnelId) ->
case persistent_term:get(macula_dist_tunnels, undefined) of
undefined -> ok;
Tunnels ->
persistent_term:put(macula_dist_tunnels, maps:remove(TunnelId, Tunnels))
end.