From 69152e86b34960b85ce6eb3a550e23b080ef60f5 Mon Sep 17 00:00:00 2001 From: Graham McIntire Date: Thu, 23 Apr 2026 14:26:19 -0500 Subject: [PATCH] refactor: split PacketProducer.handle_cast into pattern-matched heads MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit The original handle_cast had nested `if`s: outer branch on demand > 0, inner branch on buffer-over-capacity. Split into: - handle_cast with a `when demand > 0` guard for the hot path that skips the buffer entirely. - A fallback clause that computes whether the new size would overflow and delegates to buffer_packet/3. - buffer_packet/3 has two clauses — one for the drop-oldest path (with telemetry + warning) and one for the simple append path. Both share the final `maybe_activate_backpressure |> reply_noemit` step. All 19 existing PacketProducer tests pass; dialyzer stays clean. --- lib/aprsme/packet_producer.ex | 73 +++++++++++++++++++---------------- 1 file changed, 39 insertions(+), 34 deletions(-) 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}}