From 0a5527f5f6e197915ee70de53f34f5dbd774050b Mon Sep 17 00:00:00 2001 From: dajiaohuang Date: Fri, 25 Sep 2026 16:01:54 +0800 Subject: [PATCH 1/2] Validate Flow stages option --- lib/flow.ex | 6 +++++- test/flow_test.exs | 18 ++++++++++++++++++ 2 files changed, 23 insertions(+), 1 deletion(-) diff --git a/lib/flow.ex b/lib/flow.ex index 6ab10c9..e5ab69f 100644 --- a/lib/flow.ex +++ b/lib/flow.ex @@ -1538,9 +1538,13 @@ defmodule Flow do defp stages(options) do case Keyword.fetch(options, :stages) do - {:ok, _} -> + {:ok, stages} when is_integer(stages) and stages > 0 -> options + {:ok, stages} -> + raise ArgumentError, + ":stages must be a positive integer, got: #{inspect(stages)}" + :error -> stages = System.schedulers_online() [stages: stages] ++ options diff --git a/test/flow_test.exs b/test/flow_test.exs index b932e74..d7ace6e 100644 --- a/test/flow_test.exs +++ b/test/flow_test.exs @@ -115,6 +115,24 @@ defmodule FlowTest do refute_received 1 end + describe ":stages option" do + test "requires a positive integer when constructing a flow" do + constructors = [ + fn stages -> Flow.from_enumerable([1], stages: stages) end, + fn stages -> Flow.from_enumerables([[1]], stages: stages) end, + fn stages -> Flow.from_stages([self()], stages: stages) end, + fn stages -> Flow.from_specs([{Counter, 0}], stages: stages) end, + fn stages -> Flow.partition(Flow.from_enumerable([1]), stages: stages) end + ] + + for constructor <- constructors, stages <- [0, -1, :invalid] do + assert_raise ArgumentError, ~r/:stages must be a positive integer/, fn -> + constructor.(stages) + end + end + end + end + describe "errors" do test "on multiple reduce calls" do message = ~r"cannot call group_by/reduce/emit_and_reduce on a flow after another" From 8e13d324ff9eca5474ed4418df7b5fa0b3298c55 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Jos=C3=A9=20Valim?= Date: Fri, 25 Sep 2026 11:36:11 +0200 Subject: [PATCH 2/2] Apply suggestion from @josevalim --- test/flow_test.exs | 18 ------------------ 1 file changed, 18 deletions(-) diff --git a/test/flow_test.exs b/test/flow_test.exs index d7ace6e..b932e74 100644 --- a/test/flow_test.exs +++ b/test/flow_test.exs @@ -115,24 +115,6 @@ defmodule FlowTest do refute_received 1 end - describe ":stages option" do - test "requires a positive integer when constructing a flow" do - constructors = [ - fn stages -> Flow.from_enumerable([1], stages: stages) end, - fn stages -> Flow.from_enumerables([[1]], stages: stages) end, - fn stages -> Flow.from_stages([self()], stages: stages) end, - fn stages -> Flow.from_specs([{Counter, 0}], stages: stages) end, - fn stages -> Flow.partition(Flow.from_enumerable([1]), stages: stages) end - ] - - for constructor <- constructors, stages <- [0, -1, :invalid] do - assert_raise ArgumentError, ~r/:stages must be a positive integer/, fn -> - constructor.(stages) - end - end - end - end - describe "errors" do test "on multiple reduce calls" do message = ~r"cannot call group_by/reduce/emit_and_reduce on a flow after another"