Add TCP backpressure to packet pipeline

When the PacketProducer buffer crosses 80% capacity, signal Aprsme.Is
to switch its TCP socket to passive mode, leveraging TCP flow control
to slow APRS-IS intake. Resume active mode when buffer drains below 30%.

A 30-second safety valve timer auto-resumes if backpressure gets stuck.
Disconnect/reconnect/terminate paths all clear backpressure state.
This commit is contained in:
Graham McIntire 2026-02-19 15:31:38 -06:00
parent 046f97db54
commit fa89c6d3ca
No known key found for this signature in database
8 changed files with 553 additions and 19 deletions

View file

@ -68,7 +68,10 @@ config :aprsme,
# Adjust demand for smaller batches
max_demand: 300,
# Number of parallel consumers for better throughput
num_consumers: 3
num_consumers: 3,
# Backpressure water marks (ratio of max_buffer_size)
high_water_ratio: 0.8,
low_water_ratio: 0.3
]
config :error_tracker,

View file

@ -6,6 +6,7 @@ defmodule Aprsme.Is do
@aprs_timeout 60 * 1000
@keepalive_interval 20 * 1000
@backpressure_safety_valve_timeout 30_000
@type state :: %{
server: charlist() | String.t(),
@ -20,7 +21,9 @@ defmodule Aprsme.Is do
user_id: String.t(),
passcode: String.t(),
filter: String.t()
}
},
backpressure_active: boolean(),
safety_valve_timer: reference() | nil
}
def start_link(opts \\ []) do
@ -85,7 +88,9 @@ defmodule Aprsme.Is do
user_id: aprs_user_id,
passcode: aprs_passcode,
filter: default_filter
}
},
backpressure_active: false,
safety_valve_timer: nil
}
# Try to connect initially, but don't fail if it doesn't work
@ -381,15 +386,71 @@ defmodule Aprsme.Is do
{:noreply, state}
end
def handle_info({:backpressure, :activate}, %{socket: nil} = state) do
{:noreply, state}
end
def handle_info({:backpressure, :activate}, %{backpressure_active: true} = state) do
{:noreply, state}
end
def handle_info({:backpressure, :activate}, state) do
:inet.setopts(state.socket, active: false)
if state.timer, do: Process.cancel_timer(state.timer)
safety_valve_timer = Process.send_after(self(), :backpressure_safety_valve, @backpressure_safety_valve_timeout)
Logger.warning("Backpressure activated — socket set to passive mode")
{:noreply, %{state | backpressure_active: true, timer: nil, safety_valve_timer: safety_valve_timer}}
end
def handle_info({:backpressure, :deactivate}, %{backpressure_active: false} = state) do
{:noreply, state}
end
def handle_info({:backpressure, :deactivate}, %{socket: nil} = state) do
cancel_safety_valve(state)
{:noreply, %{state | backpressure_active: false, safety_valve_timer: nil}}
end
def handle_info({:backpressure, :deactivate}, state) do
:inet.setopts(state.socket, active: true)
cancel_safety_valve(state)
timer = create_timer(@aprs_timeout)
Logger.info("Backpressure deactivated — socket resumed active mode")
{:noreply, %{state | backpressure_active: false, safety_valve_timer: nil, timer: timer}}
end
def handle_info(:backpressure_safety_valve, %{backpressure_active: false} = state) do
{:noreply, state}
end
def handle_info(:backpressure_safety_valve, %{socket: nil} = state) do
{:noreply, %{state | backpressure_active: false, safety_valve_timer: nil}}
end
def handle_info(:backpressure_safety_valve, state) do
Logger.warning("Backpressure safety valve triggered — forcing resume after timeout")
:inet.setopts(state.socket, active: true)
timer = create_timer(@aprs_timeout)
{:noreply, %{state | backpressure_active: false, safety_valve_timer: nil, timer: timer}}
end
def handle_info({:tcp_closed, _socket}, state) do
Logger.warning("Socket has been closed by remote server - will reconnect")
# Cancel any existing timers
if state.timer, do: Process.cancel_timer(state.timer)
if state.keepalive_timer, do: Process.cancel_timer(state.keepalive_timer)
cancel_safety_valve(state)
# Schedule reconnect
schedule_reconnect(5000)
{:noreply, %{state | socket: nil, timer: nil, keepalive_timer: nil}}
{:noreply,
%{state | socket: nil, timer: nil, keepalive_timer: nil, backpressure_active: false, safety_valve_timer: nil}}
end
def handle_info({:tcp_error, _socket, reason}, state) do
@ -397,10 +458,13 @@ defmodule Aprsme.Is do
# Cancel any existing timers
if state.timer, do: Process.cancel_timer(state.timer)
if state.keepalive_timer, do: Process.cancel_timer(state.keepalive_timer)
cancel_safety_valve(state)
# Schedule reconnect
schedule_reconnect(5000)
{:noreply, %{state | socket: nil, timer: nil, keepalive_timer: nil}}
{:noreply,
%{state | socket: nil, timer: nil, keepalive_timer: nil, backpressure_active: false, safety_valve_timer: nil}}
end
def handle_info(:reconnect, state) do
@ -418,7 +482,16 @@ defmodule Aprsme.Is do
Logger.info("Successfully reconnected to APRS-IS")
timer = create_timer(@aprs_timeout)
keepalive_timer = create_keepalive_timer(@keepalive_interval)
{:noreply, %{state | socket: socket, timer: timer, keepalive_timer: keepalive_timer}}
{:noreply,
%{
state
| socket: socket,
timer: timer,
keepalive_timer: keepalive_timer,
backpressure_active: false,
safety_valve_timer: nil
}}
error ->
Logger.warning("Failed to login to APRS-IS: #{inspect(error)}, will retry in 10 seconds")
@ -476,6 +549,7 @@ defmodule Aprsme.Is do
# Cancel timers
if timer = Map.get(state, :timer), do: Process.cancel_timer(timer)
if timer = Map.get(state, :keepalive_timer), do: Process.cancel_timer(timer)
if timer = Map.get(state, :safety_valve_timer), do: Process.cancel_timer(timer)
# Close socket
case Map.get(state, :socket) do
@ -618,6 +692,13 @@ defmodule Aprsme.Is do
defp struct_to_map(value), do: value
defp cancel_safety_valve(%{safety_valve_timer: nil}), do: :ok
defp cancel_safety_valve(%{safety_valve_timer: timer}) do
Process.cancel_timer(timer)
:ok
end
@spec schedule_reconnect(non_neg_integer()) :: reference()
defp schedule_reconnect(delay) do
Process.send_after(self(), :reconnect, delay)

View file

@ -5,6 +5,10 @@ defmodule Aprsme.PacketProducer do
Uses `:queue` for O(1) enqueue/dequeue instead of lists, which avoids
O(n) `Enum.take` on every buffer overflow during traffic bursts.
Implements water-mark-based backpressure: when the buffer crosses the high
water mark, signals Aprsme.Is to switch its TCP socket to passive mode.
When demand drains the buffer below the low water mark, signals Is to resume.
"""
use GenStage
@ -20,15 +24,33 @@ defmodule Aprsme.PacketProducer do
@impl true
def init(opts) do
{:producer, %{demand: 0, buffer: :queue.new(), buffer_size: 0, max_buffer_size: opts[:max_buffer_size] || 1000}}
max_buffer_size = opts[:max_buffer_size] || 1000
pipeline_config = Application.get_env(:aprsme, :packet_pipeline, [])
high_ratio = Keyword.get(pipeline_config, :high_water_ratio, 0.8)
low_ratio = Keyword.get(pipeline_config, :low_water_ratio, 0.3)
{:producer,
%{
demand: 0,
buffer: :queue.new(),
buffer_size: 0,
max_buffer_size: max_buffer_size,
high_water_mark: trunc(max_buffer_size * high_ratio),
low_water_mark: trunc(max_buffer_size * low_ratio),
backpressure_active: false,
is_monitor_ref: nil
}}
end
@impl true
def handle_demand(incoming_demand, %{demand: demand} = state) do
{events, new_buffer, remaining_demand} = dispatch_events(state.buffer, demand + incoming_demand)
new_buffer_size = state.buffer_size - length(events)
{:noreply, events,
%{state | demand: remaining_demand, buffer: new_buffer, buffer_size: state.buffer_size - length(events)}}
new_state = %{state | demand: remaining_demand, buffer: new_buffer, buffer_size: new_buffer_size}
new_state = maybe_deactivate_backpressure(new_state)
{:noreply, events, new_state}
end
@impl true
@ -57,13 +79,75 @@ defmodule Aprsme.PacketProducer do
# Drop oldest (front of queue), add new to back — O(1) amortized
{_, trimmed} = :queue.out(buffer)
{:noreply, [], %{state | buffer: :queue.in(packet_data, trimmed), buffer_size: max_size}}
new_state = %{state | buffer: :queue.in(packet_data, trimmed), buffer_size: max_size}
new_state = maybe_activate_backpressure(new_state)
{:noreply, [], new_state}
else
{:noreply, [], %{state | buffer: :queue.in(packet_data, buffer), buffer_size: new_size}}
new_state = %{state | buffer: :queue.in(packet_data, buffer), buffer_size: new_size}
new_state = maybe_activate_backpressure(new_state)
{:noreply, [], new_state}
end
end
end
@impl true
def handle_info({:DOWN, ref, :process, _pid, _reason}, %{is_monitor_ref: ref} = state) do
{:noreply, [], %{state | backpressure_active: false, is_monitor_ref: nil}}
end
def handle_info({:DOWN, _ref, :process, _pid, _reason}, state) do
{:noreply, [], state}
end
def handle_info(_msg, state) do
{:noreply, [], state}
end
defp maybe_activate_backpressure(%{buffer_size: size, high_water_mark: high, backpressure_active: false} = state)
when size >= high do
case Process.whereis(Aprsme.Is) do
nil ->
state
pid ->
send(pid, {:backpressure, :activate})
ref = Process.monitor(pid)
:telemetry.execute(
[:aprsme, :packet_producer, :backpressure],
%{buffer_size: size},
%{action: :activate}
)
%{state | backpressure_active: true, is_monitor_ref: ref}
end
end
defp maybe_activate_backpressure(state), do: state
defp maybe_deactivate_backpressure(%{buffer_size: size, low_water_mark: low, backpressure_active: true} = state)
when size <= low do
if state.is_monitor_ref, do: Process.demonitor(state.is_monitor_ref, [:flush])
case Process.whereis(Aprsme.Is) do
nil ->
:ok
pid ->
send(pid, {:backpressure, :deactivate})
end
:telemetry.execute(
[:aprsme, :packet_producer, :backpressure],
%{buffer_size: size},
%{action: :deactivate}
)
%{state | backpressure_active: false, is_monitor_ref: nil}
end
defp maybe_deactivate_backpressure(state), do: state
defp dispatch_events(buffer, demand) when demand > 0 do
if :queue.is_empty(buffer) do
{[], buffer, demand}

View file

@ -108,6 +108,12 @@ defmodule AprsmeWeb.Telemetry do
reporter_options: [buckets: buckets]
),
# Backpressure Metrics
counter("aprsme.packet_producer.backpressure.count",
unit: :event,
description: "Backpressure state changes"
),
# Note: SystemMonitor and InsertOptimizer metrics removed after reverting performance optimizations
# Spatial PubSub Metrics

View file

@ -0,0 +1,96 @@
defmodule Aprsme.Is.BackpressureTest do
use ExUnit.Case, async: true
alias Aprsme.Is
# These tests call handle_info directly with constructed state,
# so they don't need a running GenServer or real socket.
defp base_state(overrides \\ []) do
defaults = %{
server: "test.aprs2.net",
port: 14_580,
socket: nil,
timer: nil,
keepalive_timer: nil,
connected_at: DateTime.utc_now(),
packet_stats: %{
total_packets: 0,
last_packet_at: nil,
packets_per_second: 0,
last_second_count: 0,
last_second_timestamp: System.system_time(:second)
},
buffer: "",
login_params: %{
user_id: "TEST",
passcode: "-1",
filter: "r/33/-96/100"
},
backpressure_active: false,
safety_valve_timer: nil
}
Enum.reduce(overrides, defaults, fn {k, v}, acc -> Map.put(acc, k, v) end)
end
describe "handle_info {:backpressure, :activate}" do
test "no-op when socket is nil" do
state = base_state()
{:noreply, new_state} = Is.handle_info({:backpressure, :activate}, state)
assert new_state.backpressure_active == false
assert new_state.safety_valve_timer == nil
end
test "no-op when already active" do
timer_ref = make_ref()
state = base_state(backpressure_active: true, safety_valve_timer: timer_ref)
{:noreply, new_state} = Is.handle_info({:backpressure, :activate}, state)
# State unchanged
assert new_state.backpressure_active == true
assert new_state.safety_valve_timer == timer_ref
end
end
describe "handle_info {:backpressure, :deactivate}" do
test "no-op when not active" do
state = base_state()
{:noreply, new_state} = Is.handle_info({:backpressure, :deactivate}, state)
assert new_state.backpressure_active == false
end
test "clears flag when socket is nil and active" do
state = base_state(backpressure_active: true, safety_valve_timer: make_ref())
{:noreply, new_state} = Is.handle_info({:backpressure, :deactivate}, state)
assert new_state.backpressure_active == false
assert new_state.safety_valve_timer == nil
end
end
describe "handle_info :backpressure_safety_valve" do
test "no-op when not active" do
state = base_state()
{:noreply, new_state} = Is.handle_info(:backpressure_safety_valve, state)
assert new_state.backpressure_active == false
end
test "clears state when socket is nil" do
state = base_state(backpressure_active: true, safety_valve_timer: make_ref())
{:noreply, new_state} = Is.handle_info(:backpressure_safety_valve, state)
assert new_state.backpressure_active == false
assert new_state.safety_valve_timer == nil
end
end
end

View file

@ -31,7 +31,9 @@ defmodule Aprsme.IsTest do
user_id: "TEST",
passcode: "-1",
filter: "r/33/-96/100"
}
},
backpressure_active: false,
safety_valve_timer: nil
}
Map.merge(base, overrides)

View file

@ -0,0 +1,193 @@
defmodule Aprsme.PacketProducerBackpressureTest do
use ExUnit.Case, async: false
alias Aprsme.PacketProducer
# We register self() as Aprsme.Is to receive backpressure messages
setup do
Process.register(self(), Aprsme.Is)
on_exit(fn ->
try do
Process.unregister(Aprsme.Is)
rescue
_ -> :ok
end
end)
:ok
end
describe "init/1 water marks" do
test "computes water marks from config ratios" do
# Config has high_water_ratio: 0.8, low_water_ratio: 0.3
{:producer, state} = PacketProducer.init(max_buffer_size: 1000)
assert state.high_water_mark == 800
assert state.low_water_mark == 300
assert state.backpressure_active == false
assert state.is_monitor_ref == nil
end
test "uses defaults when config ratios absent" do
# Temporarily clear the config
original = Application.get_env(:aprsme, :packet_pipeline)
pipeline_without_ratios = Keyword.drop(original, [:high_water_ratio, :low_water_ratio])
Application.put_env(:aprsme, :packet_pipeline, pipeline_without_ratios)
{:producer, state} = PacketProducer.init(max_buffer_size: 1000)
# Defaults: 0.8 and 0.3
assert state.high_water_mark == 800
assert state.low_water_mark == 300
Application.put_env(:aprsme, :packet_pipeline, original)
end
end
describe "backpressure activation on packet buffering" do
test "sends activate when buffer crosses high water mark" do
buffer = build_queue(799)
state = base_state(800, 300, buffer: buffer, buffer_size: 799)
# Adding one more packet brings us to 800 == high_water_mark
{:noreply, [], new_state} =
PacketProducer.handle_cast({:packet, %{sender: "TRIGGER"}}, state)
assert_received {:backpressure, :activate}
assert new_state.backpressure_active == true
assert new_state.is_monitor_ref
end
test "does not send duplicate activate when already active" do
buffer = build_queue(800)
state =
base_state(800, 300,
buffer: buffer,
buffer_size: 800,
backpressure_active: true,
is_monitor_ref: make_ref()
)
# Buffer overflows (drops oldest), but backpressure already active
{:noreply, [], new_state} =
PacketProducer.handle_cast({:packet, %{sender: "EXTRA"}}, state)
refute_received {:backpressure, :activate}
assert new_state.backpressure_active == true
end
test "skips activation when Aprsme.Is is not running" do
# Unregister so Process.whereis returns nil
Process.unregister(Aprsme.Is)
buffer = build_queue(799)
state = base_state(800, 300, buffer: buffer, buffer_size: 799)
{:noreply, [], new_state} =
PacketProducer.handle_cast({:packet, %{sender: "TRIGGER"}}, state)
# No message sent, but state should still reflect we tried
refute new_state.backpressure_active
end
end
describe "backpressure deactivation on demand" do
test "sends deactivate when demand drains buffer below low water mark" do
# Buffer at 301, demand of 2 will drain to 299 (below low_water 300)
buffer = build_queue(301)
state =
base_state(800, 300,
buffer: buffer,
buffer_size: 301,
backpressure_active: true,
is_monitor_ref: make_ref()
)
{:noreply, _events, new_state} = PacketProducer.handle_demand(2, state)
assert_received {:backpressure, :deactivate}
assert new_state.backpressure_active == false
assert new_state.is_monitor_ref == nil
end
test "does not deactivate when buffer still above low water mark" do
buffer = build_queue(305)
state =
base_state(800, 300,
buffer: buffer,
buffer_size: 305,
backpressure_active: true,
is_monitor_ref: make_ref()
)
# Demand of 2 drains to 303, still above 300
{:noreply, _events, new_state} = PacketProducer.handle_demand(2, state)
refute_received {:backpressure, :deactivate}
assert new_state.backpressure_active == true
end
end
describe "handle_info :DOWN" do
test "resets backpressure state when Is process dies" do
ref = make_ref()
state =
base_state(800, 300,
backpressure_active: true,
is_monitor_ref: ref
)
{:noreply, [], new_state} =
PacketProducer.handle_info({:DOWN, ref, :process, self(), :normal}, state)
assert new_state.backpressure_active == false
assert new_state.is_monitor_ref == nil
end
test "ignores DOWN for unrelated monitors" do
ref = make_ref()
unrelated_ref = make_ref()
state =
base_state(800, 300,
backpressure_active: true,
is_monitor_ref: ref
)
{:noreply, [], new_state} =
PacketProducer.handle_info({:DOWN, unrelated_ref, :process, self(), :normal}, state)
# Should not reset — different ref
assert new_state.backpressure_active == true
assert new_state.is_monitor_ref == ref
end
end
# Helper to build a queue of N dummy packets
defp build_queue(n) do
Enum.reduce(1..n, :queue.new(), fn i, q ->
:queue.in(%{sender: "P#{i}"}, q)
end)
end
defp base_state(high, low, overrides) do
defaults = %{
demand: 0,
buffer: :queue.new(),
buffer_size: 0,
max_buffer_size: 1000,
high_water_mark: high,
low_water_mark: low,
backpressure_active: false,
is_monitor_ref: nil
}
Enum.reduce(overrides, defaults, fn {k, v}, acc -> Map.put(acc, k, v) end)
end
end

View file

@ -25,7 +25,17 @@ defmodule Aprsme.PacketProducerTest do
describe "handle_cast {:packet, _} with demand > 0" do
test "dispatches packet immediately when demand exists" do
state = %{demand: 3, buffer: :queue.new(), buffer_size: 0, max_buffer_size: 1000}
state = %{
demand: 3,
buffer: :queue.new(),
buffer_size: 0,
max_buffer_size: 1000,
high_water_mark: 800,
low_water_mark: 300,
backpressure_active: false,
is_monitor_ref: nil
}
{:noreply, events, new_state} = PacketProducer.handle_cast({:packet, %{sender: "TEST"}}, state)
assert events == [%{sender: "TEST"}]
@ -35,7 +45,17 @@ defmodule Aprsme.PacketProducerTest do
describe "handle_cast {:packet, _} with demand == 0" do
test "buffers packet when no demand" do
state = %{demand: 0, buffer: :queue.new(), buffer_size: 0, max_buffer_size: 1000}
state = %{
demand: 0,
buffer: :queue.new(),
buffer_size: 0,
max_buffer_size: 1000,
high_water_mark: 800,
low_water_mark: 300,
backpressure_active: false,
is_monitor_ref: nil
}
{:noreply, [], new_state} = PacketProducer.handle_cast({:packet, %{sender: "TEST"}}, state)
assert new_state.buffer_size == 1
@ -49,7 +69,16 @@ defmodule Aprsme.PacketProducerTest do
:queue.in(%{sender: "PACKET#{i}"}, q)
end)
state = %{demand: 0, buffer: buffer, buffer_size: 3, max_buffer_size: 3}
state = %{
demand: 0,
buffer: buffer,
buffer_size: 3,
max_buffer_size: 3,
high_water_mark: 2,
low_water_mark: 1,
backpressure_active: false,
is_monitor_ref: nil
}
{:noreply, [], new_state} =
PacketProducer.handle_cast({:packet, %{sender: "NEW"}}, state)
@ -74,7 +103,17 @@ defmodule Aprsme.PacketProducerTest do
:queue.in(%{sender: "P#{i}"}, q)
end)
state = %{demand: 0, buffer: buffer, buffer_size: 5, max_buffer_size: 1000}
state = %{
demand: 0,
buffer: buffer,
buffer_size: 5,
max_buffer_size: 1000,
high_water_mark: 800,
low_water_mark: 300,
backpressure_active: false,
is_monitor_ref: nil
}
{:noreply, events, new_state} = PacketProducer.handle_demand(3, state)
# Should dispatch first 3 (FIFO order)
@ -87,7 +126,17 @@ defmodule Aprsme.PacketProducerTest do
test "dispatches all available when demand exceeds buffer" do
buffer = :queue.in(%{sender: "ONLY"}, :queue.new())
state = %{demand: 0, buffer: buffer, buffer_size: 1, max_buffer_size: 1000}
state = %{
demand: 0,
buffer: buffer,
buffer_size: 1,
max_buffer_size: 1000,
high_water_mark: 800,
low_water_mark: 300,
backpressure_active: false,
is_monitor_ref: nil
}
{:noreply, events, new_state} = PacketProducer.handle_demand(5, state)
assert events == [%{sender: "ONLY"}]
@ -96,14 +145,34 @@ defmodule Aprsme.PacketProducerTest do
end
test "stores demand when buffer is empty" do
state = %{demand: 0, buffer: :queue.new(), buffer_size: 0, max_buffer_size: 1000}
state = %{
demand: 0,
buffer: :queue.new(),
buffer_size: 0,
max_buffer_size: 1000,
high_water_mark: 800,
low_water_mark: 300,
backpressure_active: false,
is_monitor_ref: nil
}
{:noreply, [], new_state} = PacketProducer.handle_demand(10, state)
assert new_state.demand == 10
end
test "accumulates demand" do
state = %{demand: 5, buffer: :queue.new(), buffer_size: 0, max_buffer_size: 1000}
state = %{
demand: 5,
buffer: :queue.new(),
buffer_size: 0,
max_buffer_size: 1000,
high_water_mark: 800,
low_water_mark: 300,
backpressure_active: false,
is_monitor_ref: nil
}
{:noreply, [], new_state} = PacketProducer.handle_demand(3, state)
assert new_state.demand == 8