refactor: split PacketProducer.handle_cast into pattern-matched heads
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.
This commit is contained in:
parent
7ce011bddf
commit
69152e86b3
1 changed files with 39 additions and 34 deletions
|
|
@ -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}}
|
||||
|
|
|
|||
Loading…
Add table
Reference in a new issue