aprs.me/lib/aprsme/packet_producer.ex
Graham McIntire eb389af64b
Some checks failed
Build and Push / Build and Push Docker Image (push) Failing after 2s
reliability: bound ingress buffers, payload limits, overflow telemetry
- PacketProducer: 15s periodic buffer depth telemetry, overflow events
- MobileChannel: 500-packet ring buffer with overflow drops, payload validation
  (query max 20 chars, callsign max 20 chars, results capped at 200,
   bounds area max 1000 sq deg), rejection telemetry
- PromEx: buffer depth/max/overflow metrics, mobile overlay, payload rejected
- Tests: 2491 passed, 0 failures
- Handoff document: marked reliability #1 and #2 as done
2026-07-26 15:10:54 -05:00

191 lines
5.8 KiB
Elixir

defmodule Aprsme.PacketProducer do
@moduledoc """
GenStage producer that handles incoming APRS packets and sends them to consumers
for efficient batch processing.
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
require Logger
def start_link(opts \\ []) do
GenStage.start_link(__MODULE__, opts, name: __MODULE__)
end
def submit_packet(packet_data) do
GenStage.cast(__MODULE__, {:packet, packet_data})
end
@buffer_depth_interval_ms 15_000
@impl true
def init(opts) do
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)
buffer_depth_timer = Process.send_after(self(), :report_buffer_depth, @buffer_depth_interval_ms)
{: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,
buffer_depth_timer: buffer_depth_timer
}}
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)
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
# If a consumer is waiting, skip the buffer entirely and hand the packet
# straight through. Matches on demand > 0 directly in the head.
def handle_cast({:packet, packet_data}, %{demand: demand} = state) when demand > 0 do
{:noreply, [packet_data], %{state | demand: demand - 1}}
end
def handle_cast({:packet, packet_data}, %{buffer_size: size, max_buffer_size: max_size} = state) do
buffer_packet(size + 1 > max_size, packet_data, state)
end
# Buffer full: drop the oldest packet to make room, emit telemetry, and
# activate backpressure. O(1) amortized via :queue.
defp buffer_packet(true, packet_data, %{buffer: buffer, max_buffer_size: max_size} = state) do
dropped = state.buffer_size + 1 - max_size
Logger.warning("Packet buffer full, dropping #{dropped} oldest packet(s)",
buffer_size: max_size,
dropped: dropped
)
:telemetry.execute(
[:aprsme, :packet_producer, :buffer_overflow],
%{dropped: dropped, buffer_size: max_size},
%{}
)
:telemetry.execute(
[:aprsme, :buffer, :overflow],
%{dropped: dropped, buffer_size: max_size},
%{}
)
{_, trimmed} = :queue.out(buffer)
%{state | buffer: :queue.in(packet_data, trimmed), buffer_size: max_size}
|> maybe_activate_backpressure()
|> reply_noemit()
end
defp buffer_packet(false, packet_data, %{buffer: buffer, buffer_size: size} = state) do
%{state | buffer: :queue.in(packet_data, buffer), buffer_size: size + 1}
|> maybe_activate_backpressure()
|> reply_noemit()
end
defp reply_noemit(state), do: {:noreply, [], state}
@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(:report_buffer_depth, %{buffer_size: size, max_buffer_size: max} = state) do
:telemetry.execute(
[:aprsme, :buffer, :depth],
%{current: size, max: max},
%{}
)
timer = Process.send_after(self(), :report_buffer_depth, @buffer_depth_interval_ms)
{:noreply, [], %{state | buffer_depth_timer: timer}}
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}
else
{front, remaining} = :queue.split(min(demand, :queue.len(buffer)), buffer)
events = :queue.to_list(front)
{events, remaining, demand - length(events)}
end
end
defp dispatch_events(buffer, demand) do
{[], buffer, demand}
end
end