diff --git a/lib/gen_stage/buffer.ex b/lib/gen_stage/buffer.ex index 9852d85..d32dc37 100644 --- a/lib/gen_stage/buffer.ex +++ b/lib/gen_stage/buffer.ex @@ -122,6 +122,14 @@ defmodule GenStage.Buffer do take_count_or_until_permanent(counter, [], queue, buffer, infos) end + defp take_count_or_until_permanent(0, temps, queue, buffer, infos) + when is_reference(infos) do + {queue, buffer, perms} = + take_permanents(queue, buffer, infos, []) + + {:ok, {queue, buffer, infos}, 0, :lists.reverse(temps), :lists.reverse(perms)} + end + defp take_count_or_until_permanent(0, temps, queue, buffer, infos) do {:ok, {queue, buffer, infos}, 0, :lists.reverse(temps), []} end @@ -136,7 +144,10 @@ defmodule GenStage.Buffer do case value do {^infos, perm} -> - {:ok, {queue, buffer - 1, infos}, counter, :lists.reverse(temps), [perm]} + {queue, buffer, perms} = + take_permanents(queue, buffer - 1, infos, [perm]) + + {:ok, {queue, buffer, infos}, counter, :lists.reverse(temps), :lists.reverse(perms)} temp -> take_count_or_until_permanent(counter - 1, [temp | temps], queue, buffer - 1, infos) @@ -155,6 +166,17 @@ defmodule GenStage.Buffer do end end + defp take_permanents(queue, buffer, infos, perms) do + case :queue.peek(queue) do + {:value, {^infos, perm}} -> + {{:value, _}, queue} = :queue.out(queue) + take_permanents(queue, buffer - 1, infos, [perm | perms]) + + _ -> + {queue, buffer, perms} + end + end + ## Wheel helpers defp init_wheel(:infinity), do: make_ref() diff --git a/test/gen_stage/buffer_test.exs b/test/gen_stage/buffer_test.exs index cae390d..36105b0 100644 --- a/test/gen_stage/buffer_test.exs +++ b/test/gen_stage/buffer_test.exs @@ -132,5 +132,62 @@ defmodule GenStage.BufferTest do assert temps3 == [:temp10] assert perms3 == [:perm11] end + + test "maintains FIFO order for permanent events with infinite buffer" do + buffer = Buffer.new(:infinity) + + {buffer, _excess, _perms} = + Buffer.store_temporary(buffer, [:temp1, :temp2], :first) + + {:ok, buffer} = Buffer.store_permanent_unless_empty(buffer, :perm3) + {:ok, buffer} = Buffer.store_permanent_unless_empty(buffer, :perm4) + {:ok, buffer} = Buffer.store_permanent_unless_empty(buffer, :perm5) + + {:ok, buffer, _remaining_count, temps, perms} = + Buffer.take_count_or_until_permanent(buffer, 5) + + assert temps == [:temp1, :temp2] + assert perms == [:perm3, :perm4, :perm5] + assert Buffer.estimate_size(buffer) == 0 + end + + test "infinite buffer stops collecting permanents before next temporary event" do + buffer = Buffer.new(:infinity) + + {buffer, _, _} = Buffer.store_temporary(buffer, [:temp1], :first) + {:ok, buffer} = Buffer.store_permanent_unless_empty(buffer, :perm2) + {:ok, buffer} = Buffer.store_permanent_unless_empty(buffer, :perm3) + + {buffer, _, _} = Buffer.store_temporary(buffer, [:temp4], :first) + {:ok, buffer} = Buffer.store_permanent_unless_empty(buffer, :perm5) + + {:ok, buffer, remaining, temps, perms} = + Buffer.take_count_or_until_permanent(buffer, 5) + + assert temps == [:temp1] + assert perms == [:perm2, :perm3] + assert remaining == 4 + + {:ok, _buffer, remaining, temps, perms} = + Buffer.take_count_or_until_permanent(buffer, remaining) + + assert temps == [:temp4] + assert perms == [:perm5] + assert remaining == 3 + end + + test "infinite buffer returns permanents after the last requested temporary event" do + buffer = Buffer.new(:infinity) + {buffer, _, _} = Buffer.store_temporary(buffer, [:temp], :first) + {:ok, buffer} = Buffer.store_permanent_unless_empty(buffer, :perm1) + {:ok, buffer} = Buffer.store_permanent_unless_empty(buffer, :perm2) + + {:ok, _buffer, remaining, temps, perms} = + Buffer.take_count_or_until_permanent(buffer, 1) + + assert remaining == 0 + assert temps == [:temp] + assert perms == [:perm1, :perm2] + end end end