From fa89c6d3caaa784ce74c1a87f38148974cd7509f Mon Sep 17 00:00:00 2001 From: Graham McIntire Date: Thu, 19 Feb 2026 15:31:38 -0600 Subject: [PATCH] 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. --- config/config.exs | 5 +- lib/aprsme/is/is.ex | 91 ++++++++- lib/aprsme/packet_producer.ex | 94 ++++++++- lib/aprsme_web/telemetry.ex | 6 + test/aprsme/is/is_backpressure_test.exs | 96 +++++++++ test/aprsme/is_test.exs | 4 +- .../packet_producer_backpressure_test.exs | 193 ++++++++++++++++++ test/aprsme/packet_producer_test.exs | 83 +++++++- 8 files changed, 553 insertions(+), 19 deletions(-) create mode 100644 test/aprsme/is/is_backpressure_test.exs create mode 100644 test/aprsme/packet_producer_backpressure_test.exs diff --git a/config/config.exs b/config/config.exs index 4036942..020c038 100644 --- a/config/config.exs +++ b/config/config.exs @@ -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, diff --git a/lib/aprsme/is/is.ex b/lib/aprsme/is/is.ex index e7b4f00..2a5374c 100644 --- a/lib/aprsme/is/is.ex +++ b/lib/aprsme/is/is.ex @@ -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) diff --git a/lib/aprsme/packet_producer.ex b/lib/aprsme/packet_producer.ex index 4a8e9c1..8972466 100644 --- a/lib/aprsme/packet_producer.ex +++ b/lib/aprsme/packet_producer.ex @@ -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} diff --git a/lib/aprsme_web/telemetry.ex b/lib/aprsme_web/telemetry.ex index c1edfa8..8f4b40d 100644 --- a/lib/aprsme_web/telemetry.ex +++ b/lib/aprsme_web/telemetry.ex @@ -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 diff --git a/test/aprsme/is/is_backpressure_test.exs b/test/aprsme/is/is_backpressure_test.exs new file mode 100644 index 0000000..29db97f --- /dev/null +++ b/test/aprsme/is/is_backpressure_test.exs @@ -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 diff --git a/test/aprsme/is_test.exs b/test/aprsme/is_test.exs index 9655dab..0658138 100644 --- a/test/aprsme/is_test.exs +++ b/test/aprsme/is_test.exs @@ -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) diff --git a/test/aprsme/packet_producer_backpressure_test.exs b/test/aprsme/packet_producer_backpressure_test.exs new file mode 100644 index 0000000..a31ea4f --- /dev/null +++ b/test/aprsme/packet_producer_backpressure_test.exs @@ -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 diff --git a/test/aprsme/packet_producer_test.exs b/test/aprsme/packet_producer_test.exs index f9e1e8f..73dc0d8 100644 --- a/test/aprsme/packet_producer_test.exs +++ b/test/aprsme/packet_producer_test.exs @@ -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