Packages

macula

0.44.2
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
macula src macula_dist_system macula_dist_bridge.erl
Raw

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),
%% Socket ownership is transferred by the caller (start_supervised_bridge)
%% via gen_tcp:controlling_process AFTER start_link returns.
MonRef = erlang:monitor(process, Client),
{ok, SubRef} = subscribe_to_tunnel(Client, RecvTopic),
ReaderPid = start_reader(Client, BridgeSock, SendTopic, TunnelId, Key, Metrics),
?LOG_INFO("[dist_bridge] Started for tunnel ~s (reader ~p)", [TunnelId, ReaderPid]),
{ok, #state{bridge_sock = BridgeSock, tunnel_id = TunnelId, client = Client,
client_mon = MonRef, sub_ref = SubRef, key = Key, metrics = Metrics,
reader_pid = ReaderPid, send_topic = SendTopic, recv_topic = RecvTopic,
reconnect_count = 0},
?BRIDGE_RECV_TIMEOUT}.
handle_call(_Request, _From, State) ->
{reply, {error, unknown}, State, ?BRIDGE_RECV_TIMEOUT}.
handle_cast(_Msg, State) ->
{noreply, State, ?BRIDGE_RECV_TIMEOUT}.
%% --- 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) ->
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, ?BRIDGE_RECV_TIMEOUT};
{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, ?BRIDGE_RECV_TIMEOUT}
end;
%% --- Reader exited (linked process) ---
handle_info({'EXIT', Pid, Reason}, #state{reader_pid = Pid, tunnel_id = TunnelId} = State) ->
?LOG_INFO("[dist_bridge] Reader exited (~p) for ~s", [Reason, TunnelId]),
{stop, {reader_exit, Reason}, State};
handle_info({'EXIT', _Pid, _Reason}, State) ->
{noreply, State, ?BRIDGE_RECV_TIMEOUT};
%% --- 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);
%% --- No data timeout ---
handle_info(timeout, #state{tunnel_id = TunnelId} = State) ->
?LOG_WARNING("[dist_bridge] Timeout (no data) for ~s", [TunnelId]),
{stop, timeout, State};
handle_info(_Msg, State) ->
{noreply, State, ?BRIDGE_RECV_TIMEOUT}.
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}, ?BRIDGE_RECV_TIMEOUT};
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},
?BRIDGE_RECV_TIMEOUT}
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).
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, 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_relay_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_relay_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_relay_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.