Skip to content

Commit 6da15fb

Browse files
committed
Refactor, add specs
1 parent d626572 commit 6da15fb

13 files changed

Lines changed: 73 additions & 59 deletions

File tree

TODO.md

Lines changed: 3 additions & 1 deletion
Original file line numberDiff line numberDiff line change
@@ -13,4 +13,6 @@ Clustering:
1313

1414

1515
Refactoring
16-
- chunk transformer
16+
- chunk transformer
17+
- hanging streams after delete component
18+
- empty_gen_mix at the end of composite. delete, insert, and replace the last component

lib/composite.ex

Lines changed: 17 additions & 3 deletions
Original file line numberDiff line numberDiff line change
@@ -99,35 +99,47 @@ defmodule Strom.Composite do
9999
@spec stop(__MODULE__.t()) :: :ok
100100
def stop(%__MODULE__{} = composite), do: StartStop.stop(composite)
101101

102+
@spec delete(__MODULE__.t(), {integer(), integer()}) :: __MODULE__.t()
102103
def delete(composite, {index_from, index_to}) do
103104
GenServer.call(composite.name, {:delete, index_from, index_to})
104105
end
105106

107+
@spec delete(__MODULE__.t(), integer()) :: __MODULE__.t()
106108
def delete(composite, index) do
107109
delete(composite, {index, index})
108110
end
109111

112+
@spec insert(__MODULE__.t(), integer(), Strom.component()) :: {__MODULE__.t(), Strom.flow()}
110113
def insert(composite, index, new_component) when is_struct(new_component) do
111114
insert(composite, index, [new_component])
112115
end
113116

117+
@spec insert(__MODULE__.t(), integer(), list(Strom.component())) ::
118+
{__MODULE__.t(), Strom.flow()}
114119
def insert(composite, index, new_components) when is_list(new_components) do
115120
GenServer.call(composite.name, {:insert, index, new_components})
116121
end
117122

123+
@spec replace(__MODULE__.t(), integer(), Strom.component()) :: {__MODULE__.t(), Strom.flow()}
118124
def replace(composite, index, new_component)
119125
when is_integer(index) and is_struct(new_component) do
120126
replace(composite, {index, index}, [new_component])
121127
end
122128

129+
@spec replace(__MODULE__.t(), {integer(), integer()}, Strom.component()) ::
130+
{__MODULE__.t(), Strom.flow()}
123131
def replace(composite, {index_from, index_to}, new_component) when is_struct(new_component) do
124132
replace(composite, {index_from, index_to}, [new_component])
125133
end
126134

135+
@spec replace(__MODULE__.t(), {integer(), integer()}, list(Strom.component())) ::
136+
{__MODULE__.t(), Strom.flow()}
127137
def replace(composite, {index_from, index_to}, new_components) when is_list(new_components) do
128138
GenServer.call(composite.name, {:replace, {index_from, index_to}, new_components})
129139
end
130140

141+
@spec replace(__MODULE__.t(), integer(), list(Strom.component())) ::
142+
{__MODULE__.t(), Strom.flow()}
131143
def replace(composite, index, new_components)
132144
when is_integer(index) and is_list(new_components) do
133145
replace(composite, {index, index}, new_components)
@@ -155,9 +167,11 @@ defmodule Strom.Composite do
155167
def handle_call(
156168
{:delete, index_from, index_to},
157169
_from,
158-
%__MODULE__{components: components} = composite
170+
%__MODULE__{components: components, name: name} = composite
159171
) do
160-
{components, _deleted_components} = Manipulations.delete(components, index_from, index_to)
172+
{components, _deleted_components, %{}} =
173+
Manipulations.replace(components, index_from, index_to, [], name)
174+
161175
composite = %{composite | components: components}
162176
{:reply, composite, composite}
163177
end
@@ -168,7 +182,7 @@ defmodule Strom.Composite do
168182
%__MODULE__{components: components, name: name} = composite
169183
)
170184
when is_list(new_components) do
171-
{components, subflow} = Manipulations.insert(components, index, new_components, name)
185+
{components, [], subflow} = Manipulations.insert(components, index, new_components, name)
172186
composite = %{composite | components: components}
173187
{:reply, {composite, subflow}, composite}
174188
end

lib/composite/manipulations.ex

Lines changed: 11 additions & 37 deletions
Original file line numberDiff line numberDiff line change
@@ -5,52 +5,30 @@ defmodule Strom.Composite.Manipulations do
55
alias Strom.Composite.StartStop
66
alias Strom.GenMix
77

8-
def delete(components, index_from, index_to) do
9-
{{new_components, deleted_components}, _} =
10-
Enum.reduce(components, {{[], []}, 0}, fn component, {{acc, deleted_acc}, index} ->
11-
cond do
12-
index == index_from ->
13-
next_component = Enum.at(components, index_to + 1)
14-
input_streams = Strom.GenMix.state(component.pid).input_streams
15-
16-
{:ok, _gm_pid, new_tasks} =
17-
GenServer.call(next_component.pid, {:start_tasks, input_streams})
18-
19-
GenServer.cast(component.pid, {:transfer_tasks, new_tasks})
20-
{{acc, [component | deleted_acc]}, index + 1}
21-
22-
index > index_from and index <= index_to ->
23-
GenServer.cast(component.pid, {:gen_mix, :stopping})
24-
{{acc, [component | deleted_acc]}, index + 1}
25-
26-
true ->
27-
{{[component | acc], deleted_acc}, index + 1}
28-
end
29-
end)
30-
31-
{Enum.reverse(new_components), Enum.reverse(deleted_components)}
32-
end
33-
8+
@spec insert(list(Strom.component()), integer(), list(Strom.component()), atom()) ::
9+
{list(Strom.component()), list(Strom.component()), Strom.flow()}
3410
def insert(components, index, new_components, name) when is_list(new_components) do
3511
component_after = Enum.at(components, index)
3612
gm_after = GenMix.state(component_after.pid)
3713

3814
new_components = StartStop.start_components(new_components, name)
3915
flow = Composite.call_flow(new_components, gm_after.input_streams)
4016

41-
{:ok, _gm_pid, new_tasks} =
42-
GenServer.call(component_after.pid, {:start_tasks, Map.take(flow, component_after.inputs)})
17+
{_gm_pid, new_tasks} =
18+
GenMix.start_tasks(component_after.pid, Map.take(flow, component_after.inputs))
4319

44-
GenServer.call(component_after.pid, {:replace_tasks, new_tasks})
20+
GenMix.transfer_tasks(component_after.pid, new_tasks, :old)
4521

4622
components =
4723
components
4824
|> List.insert_at(index, new_components)
4925
|> List.flatten()
5026

51-
{components, Map.drop(flow, component_after.inputs)}
27+
{components, [], Map.drop(flow, component_after.inputs)}
5228
end
5329

30+
@spec replace(list(Strom.component()), integer(), integer(), list(Strom.component()), atom()) ::
31+
{list(Strom.component()), list(Strom.component()), Strom.flow()}
5432
def replace(components, index_from, index_to, new_components, name)
5533
when is_list(new_components) do
5634
{{new_components, deleted_components, subflow}, _} =
@@ -61,16 +39,12 @@ defmodule Strom.Composite.Manipulations do
6139
input_streams = Strom.GenMix.state(component.pid).input_streams
6240
new_components = StartStop.start_components(new_components, name)
6341
flow = Composite.call_flow(new_components, input_streams)
64-
6542
component_after = Enum.at(components, index_to + 1)
6643

67-
{:ok, _gm_pid, new_tasks} =
68-
GenServer.call(
69-
component_after.pid,
70-
{:start_tasks, Map.take(flow, component_after.inputs)}
71-
)
44+
{_gm_pid, new_tasks} =
45+
GenMix.start_tasks(component_after.pid, Map.take(flow, component_after.inputs))
7246

73-
GenServer.cast(component.pid, {:transfer_tasks, new_tasks})
47+
GenMix.transfer_tasks(component.pid, new_tasks, :all)
7448

7549
{{acc ++ Enum.reverse(new_components), [component | deleted_acc],
7650
Map.drop(flow, component_after.inputs)}, index + 1}

lib/gen_mix.ex

Lines changed: 30 additions & 12 deletions
Original file line numberDiff line numberDiff line change
@@ -33,6 +33,9 @@ defmodule Strom.GenMix do
3333
alias Strom.GenMix.Streams
3434
alias Strom.GenMix.Tasks
3535

36+
@type t() :: %__MODULE__{}
37+
38+
@spec start(__MODULE__.t()) :: __MODULE__.t()
3639
def start(%__MODULE__{process_chunk: process_chunk, opts: opts, composite: composite} = gm)
3740
when is_list(opts) do
3841
gm = %{
@@ -69,20 +72,38 @@ defmodule Strom.GenMix do
6972
end
7073

7174
@impl true
75+
@spec init(__MODULE__.t()) :: {:ok, __MODULE__.t()}
7276
def init(%__MODULE__{} = gm) do
7377
{:ok, %{gm | pid: self()}}
7478
end
7579

80+
@spec start_tasks(pid(), Strom.flow()) :: {pid(), map()}
81+
def start_tasks(gm_pid, input_streams) do
82+
GenServer.call(gm_pid, {:start_tasks, input_streams})
83+
end
84+
85+
@spec run_tasks(pid(), atom()) :: :ok
86+
def run_tasks(gm_pid, output_name) do
87+
GenServer.call(gm_pid, {:run_tasks, output_name})
88+
end
89+
7690
@spec call(map(), map()) :: map() | no_return()
7791
def call(flow, gm), do: Streams.call(flow, gm)
7892

93+
@spec state(pid()) :: __MODULE__.t()
7994
def state(pid), do: GenServer.call(pid, :state)
8095

8196
@spec stop(any()) :: any()
8297
def stop(gm) do
8398
GenServer.call(gm.pid, :stop)
8499
end
85100

101+
@spec transfer_tasks(pid(), map(), atom()) :: :ok
102+
def transfer_tasks(gm_pid, new_tasks, all_or_old) when all_or_old in [:all, :old] do
103+
GenServer.cast(gm_pid, {:transfer_tasks, new_tasks, all_or_old})
104+
end
105+
106+
@spec process_chunk(atom(), list(), Strom.flow(), any()) :: {Strom.flow(), boolean(), any()}
86107
def process_chunk(_input_stream_name, chunk, outputs, nil) do
87108
outputs
88109
|> Enum.reduce({%{}, false, nil}, fn {output_name, output_stream_fun}, {acc, any?, nil} ->
@@ -97,7 +118,7 @@ defmodule Strom.GenMix do
97118
new_tasks = Tasks.start_tasks(input_streams, gm)
98119
tasks = Map.merge(gm.tasks, new_tasks)
99120

100-
{:reply, {:ok, gm.pid, new_tasks},
121+
{:reply, {gm.pid, new_tasks},
101122
%{gm | tasks_started: true, tasks_run: false, tasks: tasks, input_streams: input_streams}}
102123
end
103124

@@ -120,16 +141,6 @@ defmodule Strom.GenMix do
120141
{:reply, :ok, %{gm | clients: clients}}
121142
end
122143

123-
def handle_call(
124-
{:replace_tasks, new_tasks},
125-
_from,
126-
%__MODULE__{tasks_started: true} = gm
127-
) do
128-
old_tasks = Map.drop(gm.tasks, Map.keys(new_tasks))
129-
run_new_tasks_and_halt_the_old_ones(new_tasks, old_tasks, gm.accs)
130-
{:reply, :ok, %{gm | tasks_run: true, tasks: new_tasks}}
131-
end
132-
133144
def handle_call(:stop, _from, %__MODULE__{} = gm) do
134145
gm = %{gm | stopping: true}
135146

@@ -151,12 +162,19 @@ defmodule Strom.GenMix do
151162
after_action(%{gm | stopping: true})
152163
end
153164

154-
def handle_cast({:transfer_tasks, new_tasks}, gm) do
165+
def handle_cast({:transfer_tasks, new_tasks, :all}, gm) do
155166
run_new_tasks_and_halt_the_old_ones(new_tasks, gm.tasks, gm.accs)
156167

157168
after_action(%{gm | stopping: true})
158169
end
159170

171+
def handle_cast({:transfer_tasks, new_tasks, :old}, gm) do
172+
old_tasks = Map.drop(gm.tasks, Map.keys(new_tasks))
173+
run_new_tasks_and_halt_the_old_ones(new_tasks, old_tasks, gm.accs)
174+
175+
after_action(%{gm | tasks_run: true, tasks: new_tasks})
176+
end
177+
160178
def handle_cast(
161179
{:put_data, {input_name, task_pid}, {new_data, new_acc}},
162180
%__MODULE__{} = gm

lib/gen_mix/streams.ex

Lines changed: 4 additions & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -1,13 +1,15 @@
11
defmodule Strom.GenMix.Streams do
22
@moduledoc "Utility module. There are functions for manipulating data in gen_mix"
33

4+
alias Strom.GenMix
5+
46
def call(flow, gm) do
57
input_streams =
68
Enum.reduce(gm.inputs, %{}, fn name, acc ->
79
Map.put(acc, name, Map.fetch!(flow, name))
810
end)
911

10-
{:ok, gm_pid, _new_tasks} = GenServer.call(gm.pid, {:start_tasks, input_streams})
12+
{gm_pid, _new_tasks} = GenMix.start_tasks(gm.pid, input_streams)
1113

1214
sub_flow = build_sub_flow(gm.outputs, gm_pid)
1315

@@ -21,7 +23,7 @@ defmodule Strom.GenMix.Streams do
2123
stream =
2224
Stream.resource(
2325
fn ->
24-
:ok = GenServer.call(gm_pid, {:run_tasks, output_name})
26+
:ok = GenMix.run_tasks(gm_pid, output_name)
2527
gm_pid
2628
end,
2729
fn gm_pid ->

lib/mixer.ex

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -54,6 +54,7 @@ defmodule Strom.Mixer do
5454
GenMix.call(flow, mixer)
5555
end
5656

57+
@spec process_chunk(atom(), list(), Strom.flow(), nil) :: {Strom.flow(), boolean(), nil}
5758
def process_chunk(_input_stream_name, chunk, outputs, nil) when map_size(outputs) == 1 do
5859
[stream_name] = Map.keys(outputs)
5960
{%{stream_name => chunk}, Enum.any?(chunk), nil}

lib/sink.ex

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -90,6 +90,7 @@ defmodule Strom.Sink do
9090
end)
9191
end
9292

93+
@spec process_chunk(atom(), list(), Strom.flow(), any()) :: {Strom.flow(), false, nil}
9394
def process_chunk(input_stream_name, _chunk, _outputs, nil) do
9495
{%{input_stream_name => []}, false, nil}
9596
end

lib/splitter.ex

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -70,6 +70,7 @@ defmodule Strom.Splitter do
7070
GenMix.call(flow, splitter)
7171
end
7272

73+
@spec process_chunk(atom(), list(), Strom.flow(), nil) :: {Strom.flow(), boolean(), nil}
7374
def process_chunk(_input_stream_name, chunk, outputs, nil) do
7475
outputs
7576
|> Enum.reduce({%{}, false, nil}, fn {stream_name, fun}, {acc, any?, nil} ->

lib/strom.ex

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -3,6 +3,7 @@ defmodule Strom do
33

44
@type event() :: any()
55
@type stream_name() :: any()
6+
@type component() :: struct()
67
@type stream() :: Enumerable.t(event())
78
@type flow() :: %{optional(stream_name()) => stream()}
89
end

lib/transformer.ex

Lines changed: 1 addition & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -88,6 +88,7 @@ defmodule Strom.Transformer do
8888
GenMix.call(flow, transformer)
8989
end
9090

91+
@spec process_chunk(atom(), list(), Strom.flow(), any()) :: {Strom.flow(), boolean(), any()}
9192
def process_chunk(input_stream_name, chunk, outputs, acc) do
9293
output_function = Map.get(outputs, input_stream_name)
9394

0 commit comments

Comments
 (0)