Packages
Implementation of stable OpenTelemetry signals
Retired package: Release invalid - Package has a breaking bug
Current section
Files
Jump to
Current section
Files
src/otel_simple_processor.erl
%%%------------------------------------------------------------------------
%% Copyright 2021, OpenTelemetry Authors
%% Licensed under the Apache License, Version 2.0 (the "License");
%% you may not use this file except in compliance with the License.
%% You may obtain a copy of the License at
%%
%% http://www.apache.org/licenses/LICENSE-2.0
%%
%% Unless required by applicable law or agreed to in writing, software
%% distributed under the License is distributed on an "AS IS" BASIS,
%% WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
%% See the License for the specific language governing permissions and
%% limitations under the License.
%%
%% @doc This Span Processor synchronously exports each ended Span.
%%
%% Use this processor if ending a Span should block until it has been
%% exported. This is useful for cases like a serverless environment where
%% the application will possibly be suspended after handling a request.
%%
%% @end
%%%-----------------------------------------------------------------------
-module(otel_simple_processor).
-behaviour(gen_statem).
-behaviour(otel_span_processor).
-export([start_link/1,
on_start/3,
on_end/2,
force_flush/1,
report_cb/1,
%% deprecated
set_exporter/1,
set_exporter/2,
set_exporter/3]).
-export([init/1,
callback_mode/0,
idle/3,
exporting/3,
terminate/3]).
%% uncomment when OTP-23 becomes the minimum required version
%% -deprecated({set_exporter, 1, "set through the otel_tracer_provider instead"}).
%% -deprecated({set_exporter, 2, "set through the otel_tracer_provider instead"}).
%% -deprecated({set_exporter, 3, "set through the otel_tracer_provider instead"}).
-eqwalizer({nowarn_function, on_end/2}).
-include_lib("opentelemetry_api/include/opentelemetry.hrl").
-include_lib("kernel/include/logger.hrl").
-include("otel_span.hrl").
-record(data, {exporter :: {module(), term()} | undefined,
exporter_config :: {module(), term()} | undefined | none,
current_from :: gen_statem:from() | undefined,
resource :: otel_resource:t(),
handed_off_table :: atom() | undefined,
runner_pid :: pid() | undefined,
exporting_timeout_ms :: integer()}).
-define(DEFAULT_EXPORTER_TIMEOUT_MS, timer:minutes(5)).
-define(NAME_TO_ATOM(Name, Unique), list_to_atom(lists:concat([Name, "_", Unique]))).
%% require a unique name to distinguish multiple simple processors while
%% still having a single name, instead of a possibly changing pid, to
%% communicate with the processor
%% @doc Starts a Simple Span Processor.
%% @end
-spec start_link(#{name := atom() | list()}) -> {ok, pid(), map()}.
start_link(Config=#{name := Name}) ->
RegisterName = ?NAME_TO_ATOM(?MODULE, Name),
Config1 = Config#{reg_name => RegisterName},
{ok, Pid} = gen_statem:start_link({local, RegisterName}, ?MODULE, [Config1], []),
{ok, Pid, Config1}.
%% @deprecated Please use {@link otel_tracer_provider}
set_exporter(Exporter) ->
set_exporter(global, Exporter, []).
%% @deprecated Please use {@link otel_tracer_provider}
-spec set_exporter(module(), term()) -> ok.
set_exporter(Exporter, Options) ->
gen_statem:call(?REG_NAME(global), {set_exporter, {Exporter, Options}}).
%% @deprecated Please use {@link otel_tracer_provider}
-spec set_exporter(atom(), module(), term()) -> ok.
set_exporter(Name, Exporter, Options) ->
gen_statem:call(?REG_NAME(Name), {set_exporter, {Exporter, Options}}).
%% @private
-spec on_start(otel_ctx:t(), opentelemetry:span(), otel_span_processor:processor_config())
-> opentelemetry:span().
on_start(_Ctx, Span, _) ->
Span.
%% @private
-spec on_end(opentelemetry:span(), otel_span_processor:processor_config())
-> true | dropped | {error, invalid_span} | {error, no_export_buffer}.
on_end(#span{trace_flags=TraceFlags}, _) when not(?IS_SAMPLED(TraceFlags)) ->
dropped;
on_end(Span=#span{}, #{reg_name := RegName}) ->
gen_statem:call(RegName, {export, Span});
on_end(_Span, _) ->
{error, invalid_span}.
%% @private
-spec force_flush(#{reg_name := gen_statem:server_ref()}) -> ok.
force_flush(#{reg_name := RegName}) ->
gen_statem:cast(RegName, force_flush).
%% @private
init([Args]) ->
process_flag(trap_exit, true),
ExportingTimeout = maps:get(exporting_timeout_ms, Args, ?DEFAULT_EXPORTER_TIMEOUT_MS),
%% TODO: this should be passed in from the tracer server
Resource = case maps:find(resource, Args) of
{ok, R} ->
R;
error ->
otel_resource_detector:get_resource()
end,
%% Resource = otel_tracer_provider:resource(),
{ok, idle, #data{exporter=undefined,
exporter_config=maps:get(exporter, Args, none),
resource = Resource,
handed_off_table=undefined,
exporting_timeout_ms=ExportingTimeout},
[{next_event, internal, init_exporter}]}.
%% @private
callback_mode() ->
state_functions.
%% @private
idle({call, From}, {export, _Span}, #data{exporter=undefined}) ->
{keep_state_and_data, [{reply, From, dropped}]};
idle({call, From}, {export, Span}, Data) ->
{next_state, exporting, Data, [{next_event, internal, {export, From, Span}}]};
idle(EventType, Event, Data) ->
handle_event_(idle, EventType, Event, Data).
%% @private
exporting({call, _From}, {export, _}, _) ->
{keep_state_and_data, [postpone]};
exporting(internal, {export, From, Span}, Data=#data{exporting_timeout_ms=ExportingTimeout}) ->
{OldTableName, RunnerPid} = export_span(Span, Data),
{keep_state, Data#data{runner_pid=RunnerPid,
current_from=From,
handed_off_table=OldTableName},
[{state_timeout, ExportingTimeout, exporting_timeout}]};
exporting(state_timeout, exporting_timeout, Data=#data{current_from=From,
handed_off_table=_ExportingTable}) ->
%% kill current exporting process because it is taking too long
%% which deletes the exporting table, so create a new one and
%% repeat the state to force another span exporting immediately
Data1 = kill_runner(Data),
{next_state, idle, Data1, [{reply, From, {error, timeout}}]};
%% important to verify runner_pid and FromPid are the same in case it was sent
%% after kill_runner was called but before it had done the unlink
exporting(info, {'EXIT', FromPid, _}, Data=#data{runner_pid=FromPid}) ->
complete_exporting(Data);
%% important to verify runner_pid and FromPid are the same in case it was sent
%% after kill_runner was called but before it had done the unlink
exporting(info, {completed, FromPid}, Data=#data{runner_pid=FromPid}) ->
complete_exporting(Data);
exporting(EventType, Event, Data) ->
handle_event_(exporting, EventType, Event, Data).
handle_event_(_, internal, init_exporter, Data=#data{exporter=undefined,
exporter_config=ExporterConfig}) ->
Exporter = otel_exporter:init(ExporterConfig),
{keep_state, Data#data{exporter=Exporter}};
handle_event_(_, {call, From}, {set_exporter, Exporter}, Data=#data{exporter=OldExporter}) ->
otel_exporter:shutdown(OldExporter),
{keep_state, Data#data{exporter=otel_exporter:init(Exporter)}, [{reply, From, ok}]};
handle_event_(_, _, _, _) ->
keep_state_and_data.
%% @private
terminate(_, _, _Data) ->
ok.
%%
complete_exporting(Data=#data{current_from=From,
handed_off_table=ExportingTable})
when ExportingTable =/= undefined ->
{next_state, idle, Data#data{current_from=undefined,
runner_pid=undefined,
handed_off_table=undefined},
[{reply, From, ok}]}.
kill_runner(Data=#data{runner_pid=RunnerPid}) when RunnerPid =/= undefined ->
Mon = erlang:monitor(process, RunnerPid),
erlang:unlink(RunnerPid),
erlang:exit(RunnerPid, kill),
%% Wait for the runner process termination to be sure that
%% the export table is destroyed and can be safely recreated
receive
{'DOWN', Mon, process, RunnerPid, _} ->
Data#data{runner_pid=undefined, handed_off_table=undefined}
end.
new_export_table(Name) ->
ets:new(Name, [public,
{write_concurrency, true},
duplicate_bag,
%% OpenTelemetry exporter protos group by the
%% instrumentation_scope. So using instrumentation_scope
%% as the key means we can easily lookup all spans for
%% for each instrumentation_scope and export together.
{keypos, #span.instrumentation_scope}]).
export_span(Span, #data{exporter=Exporter,
resource=Resource}) ->
Table = new_export_table(otel_simple_processor_table),
_ = ets:insert(Table, Span),
Self = self(),
RunnerPid = erlang:spawn_link(fun() -> send_spans(Self, Resource, Exporter) end),
ets:give_away(Table, RunnerPid, export),
{Table, RunnerPid}.
%% Additional benefit of using a separate process is calls to `register` won't
%% timeout if the actual exporting takes longer than the call timeout
send_spans(FromPid, Resource, Exporter) ->
receive
{'ETS-TRANSFER', Table, FromPid, export} ->
export(Exporter, Resource, Table),
ets:delete(Table),
completed(FromPid)
end.
completed(FromPid) ->
FromPid ! {completed, self()}.
export(undefined, _, _) ->
true;
export(Exporter, Resource, SpansTid) ->
%% don't let a exporter exception crash us
%% and return true if exporter failed
try
otel_exporter_traces:export(Exporter, SpansTid, Resource) =:= failed_not_retryable
catch
Kind:Reason:StackTrace ->
?LOG_INFO(#{source => exporter,
during => export,
kind => Kind,
reason => Reason,
exporter => Exporter,
stacktrace => StackTrace}, #{report_cb => fun ?MODULE:report_cb/1}),
true
end.
%% logger format functions
%% @private
report_cb(#{source := exporter,
during := export,
kind := Kind,
reason := Reason,
exporter := {ExporterModule, _},
stacktrace := StackTrace}) ->
{"exporter threw exception: exporter=~p ~ts",
[ExporterModule, otel_utils:format_exception(Kind, Reason, StackTrace)]}.