Current section

8 Versions

Jump to

Compare versions

5 files changed
+187 additions
-107 deletions
  @@ -5,9 +5,9 @@
5 5 {<<"files">>,
6 6 [<<"lib">>,<<"lib/orion_collector">>,
7 7 <<"lib/orion_collector/application.ex">>,
8 - <<"lib/orion_collector/tracer.ex">>,<<"lib/orion_collector.ex">>,
9 - <<".formatter.exs">>,<<"mix.exs">>,<<"README.md">>,<<"LICENSE">>,
10 - <<"CHANGELOG.md">>]}.
8 + <<"lib/orion_collector/tracer.ex">>,<<"lib/orion_collector/aggregator.ex">>,
9 + <<"lib/orion_collector.ex">>,<<".formatter.exs">>,<<"mix.exs">>,
10 + <<"README.md">>,<<"LICENSE">>,<<"CHANGELOG.md">>]}.
11 11 {<<"licenses">>,[<<"Apache-2.0">>]}.
12 12 {<<"links">>,
13 13 [{<<"github">>,<<"https://github.com/LivewareProblems/orion_collector">>}]}.
  @@ -23,4 +23,4 @@
23 23 {<<"optional">>,false},
24 24 {<<"repository">>,<<"hexpm">>},
25 25 {<<"requirement">>,<<"~> 1.6">>}]]}.
26 - {<<"version">>,<<"1.0.1">>}.
26 + {<<"version">>,<<"1.1.0">>}.
  @@ -0,0 +1,136 @@
1 + defmodule OrionCollector.Aggregator do
2 + use GenServer
3 + import Ex2ms
4 +
5 + def start_agg(mfa, pid) do
6 + spec_args = {OrionCollector.Aggregator, [mfa, pid]}
7 + DynamicSupervisor.start_child(OrionCollector.AggregatorSupervisor, spec_args)
8 + end
9 +
10 + def mfa_to_name(mfa) do
11 + {:via, Registry, {OrionCollector.Aggregator.Registry, mfa}}
12 + end
13 +
14 + def stop(pid) do
15 + GenServer.stop(pid, :normal, 5_000)
16 + end
17 +
18 + def start_link(init = [mfa, pid]) do
19 + GenServer.start_link(__MODULE__, init,
20 + name: {:via, Registry, {OrionCollector.Aggregator.Registry, mfa, pid}}
21 + )
22 + end
23 +
24 + def child_spec(args) do
25 + %{
26 + id: {OrionCollector.Aggregator, args},
27 + start: {OrionCollector.Aggregator, :start_link, [args]},
28 + restart: :transient
29 + }
30 + end
31 +
32 + # --PRIVATE--
33 + @impl true
34 + def init([mfa, pid]) do
35 + mon_ref = Process.monitor(pid)
36 +
37 + :erlang.trace_pattern(mfa, match_spec(), [:local])
38 +
39 + Process.send_after(self(), :send_data, 500)
40 +
41 + initial_state = %{
42 + mfa: mfa,
43 + liveview_pid: pid,
44 + call_depth: %{},
45 + time_stored: %{},
46 + ddsketch: DogSketch.SimpleDog.new(),
47 + ref_mon: mon_ref
48 + }
49 +
50 + {:ok, initial_state}
51 + end
52 +
53 + @impl true
54 + def handle_cast(
55 + {:trace_ts, trace_pid, :call, _mfa, start_time},
56 + state = %{call_depth: cd_map, time_stored: time_stored_map}
57 + ) do
58 + cd = Map.get(cd_map, trace_pid, 0)
59 + new_cd_map = Map.put(cd_map, trace_pid, cd + 1)
60 +
61 + new_ts_map =
62 + if cd == 0 do
63 + Map.put(time_stored_map, {trace_pid, cd + 1}, start_time)
64 + else
65 + time_stored_map
66 + end
67 +
68 + new_state =
69 + state
70 + |> Map.put(:call_depth, new_cd_map)
71 + |> Map.put(:time_stored, new_ts_map)
72 +
73 + {:noreply, new_state}
74 + end
75 +
76 + @accepted_return_tags [:return_from, :exception_from]
77 +
78 + @impl true
79 + def handle_cast(
80 + {:trace_ts, trace_pid, return_tag, _mfa, _TraceTerm, end_time},
81 + state = %{call_depth: cd_map, time_stored: time_stored_map, ddsketch: ddsketch}
82 + )
83 + when return_tag in @accepted_return_tags do
84 + case Map.get(cd_map, trace_pid, 0) do
85 + 0 ->
86 + {:noreply, state}
87 +
88 + 1 ->
89 + new_cd_map = Map.delete(cd_map, trace_pid)
90 + {start_time, new_ts_map} = Map.pop(time_stored_map, {trace_pid, 1})
91 +
92 + call_time_micro = :timer.now_diff(end_time, start_time)
93 + new_sketch = DogSketch.SimpleDog.insert(ddsketch, call_time_micro / 1_000)
94 +
95 + new_state =
96 + state
97 + |> Map.put(:call_depth, new_cd_map)
98 + |> Map.put(:time_stored, new_ts_map)
99 + |> Map.put(:ddsketch, new_sketch)
100 +
101 + {:noreply, new_state}
102 +
103 + cd when cd > 1 ->
104 + new_cd_map = Map.put(cd_map, trace_pid, cd - 1)
105 +
106 + {:noreply, Map.put(state, :call_depth, new_cd_map)}
107 + end
108 + end
109 +
110 + @impl true
111 + def handle_info(:send_data, %{ddsketch: ddsketch, liveview_pid: liveview_pid} = state) do
112 + new_sketch = DogSketch.SimpleDog.new()
113 +
114 + if ddsketch != new_sketch do
115 + send(liveview_pid, {:ddsketch, ddsketch})
116 + end
117 +
118 + Process.send_after(self(), :send_data, 500)
119 +
120 + {:noreply, Map.put(state, :ddsketch, new_sketch)}
121 + end
122 +
123 + @impl true
124 + def handle_info({:DOWN, ref, :process, _object, _reason}, %{ref_mon: ref, mfa: mfa} = state) do
125 + :erlang.trace_pattern(mfa, false, [])
126 + {:stop, :normal, state}
127 + end
128 +
129 + defp match_spec() do
130 + fun do
131 + _ ->
132 + return_trace()
133 + exception_trace()
134 + end
135 + end
136 + end
  @@ -8,7 +8,9 @@ defmodule OrionCollector.Application do
8 8 @impl true
9 9 def start(_type, _args) do
10 10 children = [
11 - {DynamicSupervisor, strategy: :one_for_one, name: OrionCollector.TracerSupervisor}
11 + {Registry, keys: :unique, name: OrionCollector.Aggregator.Registry},
12 + OrionCollector.Tracer,
13 + {DynamicSupervisor, strategy: :one_for_one, name: OrionCollector.AggregatorSupervisor}
12 14
13 15 # Starts a worker by calling: OrionCollector.Worker.start_link(arg)
14 16 # {OrionCollector.Worker, arg}
  @@ -16,7 +18,7 @@ defmodule OrionCollector.Application do
16 18
17 19 # See https://hexdocs.pm/elixir/Supervisor.html
18 20 # for other strategies and supported options
19 - opts = [strategy: :one_for_one, name: OrionCollector.Supervisor]
21 + opts = [strategy: :rest_for_one, name: OrionCollector.Supervisor]
20 22 Supervisor.start_link(children, opts)
21 23 end
22 24 end
  @@ -1,7 +1,13 @@
1 1 defmodule OrionCollector.Tracer do
2 2 use GenServer
3 - import Ex2ms
4 3
4 + alias OrionCollector.Aggregator
5 +
6 + @moduledoc """
7 + This is the process that collect the Trace, then dispatch them per aggregator
8 +
9 + There is one Tracer per node, but one Aggregator per MFA being traced.
10 + """
5 11 def start_all_node_tracers(mfa, self, start_status \\ :running) do
6 12 if self do
7 13 OrionCollector.Tracer.start_tracer(mfa, self(), start_status)
  @@ -22,31 +28,30 @@ defmodule OrionCollector.Tracer do
22 28 :ok
23 29 end
24 30
25 - def pause_trace(mfa, self) do
31 + def pause_trace(self) do
26 32 if self do
27 - :erlang.trace_pattern(mfa, false, [:local])
33 + OrionCollector.Tracer.change_status(:pause)
28 34 end
29 35
30 - :erpc.multicall(list_nodes(), :erlang, :trace_pattern, [mfa, false, [:local]], 5_000)
36 + :erpc.multicall(list_nodes(), OrionCollector.Tracer, :change_status, [:pause], 5_000)
31 37 end
32 38
33 - def restart_trace(mfa, self) do
39 + def restart_trace(self) do
34 40 if self do
35 - :erlang.trace_pattern(mfa, match_spec(), [:local])
41 + OrionCollector.Tracer.change_status(:start)
36 42 end
37 43
38 - :erpc.multicall(
39 - list_nodes(),
40 - :erlang,
41 - :trace_pattern,
42 - [mfa, match_spec(), [:local]],
43 - 5_000
44 - )
44 + :erpc.multicall(list_nodes(), OrionCollector.Tracer, :change_status, [:start], 5_000)
45 45 end
46 46
47 47 def start_tracer(mfa, pid, start_status) do
48 - spec_args = {OrionCollector.Tracer, [mfa, pid, start_status]}
49 - DynamicSupervisor.start_child(OrionCollector.TracerSupervisor, spec_args)
48 + Aggregator.start_agg(mfa, pid)
49 + change_status(start_status)
50 + end
51 +
52 + @spec change_status(any) :: any
53 + def change_status(status) do
54 + GenServer.call(__MODULE__, status)
50 55 end
51 56
52 57 def stop(pid) do
  @@ -54,119 +59,56 @@ defmodule OrionCollector.Tracer do
54 59 end
55 60
56 61 def start_link(init \\ []) do
57 - GenServer.start_link(__MODULE__, init)
62 + GenServer.start_link(__MODULE__, init, name: __MODULE__)
58 63 end
59 64
60 65 def child_spec(args) do
61 66 %{
62 - id: {OrionCollector.Tracer, args},
67 + id: OrionCollector.Tracer,
63 68 start: {OrionCollector.Tracer, :start_link, [args]},
64 - restart: :transient
69 + restart: :permanent
65 70 }
66 71 end
67 72
68 73 # --PRIVATE--
69 74 @impl true
70 - def init([mfa, pid, start_status]) do
71 - mon_ref = Process.monitor(pid)
75 + def init(_start) do
76 + running_trace(false)
72 77
73 - if start_status == :running do
74 - :erlang.trace_pattern(mfa, match_spec(), [:local])
75 - end
76 -
77 - :erlang.trace(:all, true, [:call, :arity, :timestamp])
78 -
79 - Process.send_after(self(), :send_data, 500)
80 -
81 - initial_state = %{
82 - mfa: mfa,
83 - liveview_pid: pid,
84 - call_depth: %{},
85 - time_stored: %{},
86 - ddsketch: DogSketch.SimpleDog.new(),
87 - ref_mon: mon_ref
88 - }
78 + initial_state = %{running_status: :pause}
89 79
90 80 {:ok, initial_state}
91 81 end
92 82
83 + @impl true
84 + def handle_call(start_status, _from, state) do
85 + running_trace(start_status == :start)
86 + {:reply, :ok, Map.put(state, :running_status, start_status)}
87 + end
88 +
93 89 @impl true
94 90 def handle_info(
95 - {:trace_ts, trace_pid, :call, _mfa, start_time},
96 - state = %{call_depth: cd_map, time_stored: time_stored_map}
91 + {:trace_ts, _trace_pid, :call, mfa, _start_time} = trace_msg,
92 + state
97 93 ) do
98 - cd = Map.get(cd_map, trace_pid, 0)
99 - new_cd_map = Map.put(cd_map, trace_pid, cd + 1)
100 -
101 - new_ts_map =
102 - if cd == 0 do
103 - Map.put(time_stored_map, {trace_pid, cd + 1}, start_time)
104 - else
105 - time_stored_map
106 - end
107 -
108 - new_state =
109 - state
110 - |> Map.put(:call_depth, new_cd_map)
111 - |> Map.put(:time_stored, new_ts_map)
112 -
113 - {:noreply, new_state}
94 + GenServer.cast(Aggregator.mfa_to_name(mfa), trace_msg)
95 + {:noreply, state}
114 96 end
115 97
116 98 @accepted_return_tags [:return_from, :exception_from]
117 99
118 100 @impl true
119 101 def handle_info(
120 - {:trace_ts, trace_pid, return_tag, _mfa, _TraceTerm, end_time},
121 - state = %{call_depth: cd_map, time_stored: time_stored_map, ddsketch: ddsketch}
102 + {:trace_ts, _trace_pid, return_tag, mfa, _TraceTerm, _end_time} = trace_msg,
103 + state
122 104 )
123 105 when return_tag in @accepted_return_tags do
124 - case Map.get(cd_map, trace_pid, 0) do
125 - 0 ->
126 - {:noreply, state}
127 -
128 - 1 ->
129 - new_cd_map = Map.delete(cd_map, trace_pid)
130 - {start_time, new_ts_map} = Map.pop(time_stored_map, {trace_pid, 1})
131 -
132 - call_time_micro = :timer.now_diff(end_time, start_time)
133 - new_sketch = DogSketch.SimpleDog.insert(ddsketch, call_time_micro / 1_000)
134 -
135 - new_state =
136 - state
137 - |> Map.put(:call_depth, new_cd_map)
138 - |> Map.put(:time_stored, new_ts_map)
139 - |> Map.put(:ddsketch, new_sketch)
140 -
141 - {:noreply, new_state}
142 -
143 - cd when cd > 1 ->
144 - new_cd_map = Map.put(cd_map, trace_pid, cd - 1)
145 -
146 - {:noreply, Map.put(state, :call_depth, new_cd_map)}
147 - end
106 + GenServer.cast(Aggregator.mfa_to_name(mfa), trace_msg)
107 + {:noreply, state}
148 108 end
149 109
150 - @impl true
151 - def handle_info(:send_data, %{ddsketch: ddsketch, liveview_pid: liveview_pid} = state) do
152 - send(liveview_pid, {:ddsketch, ddsketch})
153 - Process.send_after(self(), :send_data, 500)
154 -
155 - {:noreply, Map.put(state, :ddsketch, DogSketch.SimpleDog.new())}
156 - end
157 -
158 - @impl true
159 - def handle_info({:DOWN, ref, :process, _object, _reason}, %{ref_mon: ref, mfa: mfa} = state) do
160 - :erlang.trace_pattern(mfa, false, [])
161 - {:stop, :normal, state}
162 - end
163 -
164 - defp match_spec() do
165 - fun do
166 - _ ->
167 - return_trace()
168 - exception_trace()
169 - end
110 + defp running_trace(bool) do
111 + :erlang.trace(:all, bool, [:call, :arity, :timestamp])
170 112 end
171 113
172 114 defp list_nodes() do
  @@ -1,7 +1,7 @@
1 1 defmodule OrionCollector.MixProject do
2 2 use Mix.Project
3 3
4 - @version "1.0.1"
4 + @version "1.1.0"
5 5
6 6 def project do
7 7 [