diff --git a/lib/aprsme/packet_producer.ex b/lib/aprsme/packet_producer.ex index 8972466..4f8619b 100644 --- a/lib/aprsme/packet_producer.ex +++ b/lib/aprsme/packet_producer.ex @@ -54,42 +54,47 @@ defmodule Aprsme.PacketProducer do end @impl true - def handle_cast( - {:packet, packet_data}, - %{demand: demand, buffer: buffer, buffer_size: size, max_buffer_size: max_size} = state - ) do - if demand > 0 do - {:noreply, [packet_data], %{state | demand: demand - 1}} - else - new_size = size + 1 - - if new_size > max_size do - dropped = new_size - 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}, - %{} - ) - - # Drop oldest (front of queue), add new to back — O(1) amortized - {_, trimmed} = :queue.out(buffer) - new_state = %{state | buffer: :queue.in(packet_data, trimmed), buffer_size: max_size} - new_state = maybe_activate_backpressure(new_state) - {:noreply, [], new_state} - else - 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 + # 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}, + %{} + ) + + {_, 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}}