Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
18 changes: 18 additions & 0 deletions README.md
Original file line number Diff line number Diff line change
Expand Up @@ -177,6 +177,24 @@ defmodule MyExternalService do
end
```

When an "Other" Transaction is started automatically but the trace headers arrive some other way, such as in Oban job metadata or a message attribute, connect it to the trace with `NewRelic.accept_distributed_trace_headers/1`:

```elixir
defmodule MyWorker do
use Oban.Worker

def enqueue(args) do
Oban.insert(new(args, meta: %{dt_headers: NewRelic.distributed_trace_headers(:other)}))
end

@impl Oban.Worker
def perform(%Oban.Job{meta: meta}) do
NewRelic.accept_distributed_trace_headers(meta["dt_headers"])
# ...
end
end
```

#### Mix Tasks

`NewRelic.Instrumented.Mix.Task` To enable the agent and record an Other Transaction during a `Mix.Task`, simply `use NewRelic.Instrumented.Mix.Task`. This will ensure the agent is properly started, records a Transaction, and is shut down.
Expand Down
2 changes: 1 addition & 1 deletion examples/apps/oban_example/README.md
Original file line number Diff line number Diff line change
@@ -1,3 +1,3 @@
# ObanExample

An example app demonstrating auto-instrumentation of Oban
An example app demonstrating auto-instrumentation of Oban, including connecting a job to the Distributed Trace that enqueued it with `NewRelic.accept_distributed_trace_headers/1`.
4 changes: 3 additions & 1 deletion examples/apps/oban_example/lib/oban_example/worker.ex
Original file line number Diff line number Diff line change
Expand Up @@ -6,7 +6,9 @@ defmodule ObanExample.Worker do
{:error, message}
end

def perform(%Oban.Job{args: _args}) do
def perform(%Oban.Job{meta: meta}) do
NewRelic.accept_distributed_trace_headers(meta["dt_headers"])

Process.sleep(15 + :rand.uniform(50))
:ok
end
Expand Down
26 changes: 26 additions & 0 deletions examples/apps/oban_example/test/oban_example_test.exs
Original file line number Diff line number Diff line change
Expand Up @@ -37,6 +37,32 @@ defmodule ObanExampleTest do
assert event[:"oban.job.tags"] == "foo,bar"
end

test "connects a job to the Distributed Trace that enqueued it" do
TestHelper.restart_harvest_cycle(Collector.TransactionEvent.HarvestCycle)

dt_headers =
Task.async(fn ->
NewRelic.start_transaction("Test", "Origin")
NewRelic.distributed_trace_headers(:other)
end)
|> Task.await()

ObanExample.Worker.new(%{some: "args"}, meta: %{dt_headers: dt_headers})
|> Oban.insert()

events = TestHelper.gather_harvest(Collector.TransactionEvent.Harvester)

origin = TestHelper.find_event(events, "OtherTransaction/Test/Origin")

job =
TestHelper.find_event(events, "OtherTransaction/Oban/default/ObanExample.Worker/perform")

assert job[:traceId] == origin[:traceId]
assert job[:parentId] == origin[:guid]
assert job[:"parent.type"] == "App"
assert job[:"parent.transportType"] == "Other"
end

test "instruments a failed job" do
TestHelper.restart_harvest_cycle(Collector.Metric.HarvestCycle)
TestHelper.restart_harvest_cycle(Collector.TransactionEvent.HarvestCycle)
Expand Down
33 changes: 33 additions & 0 deletions lib/new_relic.ex
Original file line number Diff line number Diff line change
Expand Up @@ -277,11 +277,44 @@ defmodule NewRelic do

* Call `distributed_trace_headers` immediately before making the
request since calling the function marks the "start" time of the request.
* Returns an empty list (or map) when the agent is not enabled or
when called outside of a Transaction.
"""
@spec distributed_trace_headers(:http) :: [{key :: String.t(), value :: String.t()}]
@spec distributed_trace_headers(:other) :: map()
defdelegate distributed_trace_headers(type), to: NewRelic.DistributedTrace

@doc """
Connect the current "Other" Transaction to an existing Distributed Trace.

Use this when an "Other" Transaction is started automatically (for example by
the Oban instrumentation) but the trace headers arrive some other way, such as
in job metadata or a message attribute. The headers can be W3C "traceparent"
and "tracestate" headers or another New Relic agent's "newrelic" header.

```elixir
# When enqueueing
Oban.insert(MyWorker.new(args, meta: %{dt_headers: NewRelic.distributed_trace_headers(:other)}))

# When performing
def perform(%Oban.Job{meta: meta}) do
NewRelic.accept_distributed_trace_headers(meta["dt_headers"])
# ...
end
```

## Notes

* Web Transactions read inbound headers automatically, so this is only
needed for "Other" Transactions.
* Call this as early as possible in the Transaction so all Spans are linked.
* Ignored when called outside of a Transaction, when the headers can't be
decoded, when the Transaction already accepted inbound headers, or when the
agent is disabled.
"""
@spec accept_distributed_trace_headers(headers :: map()) :: :ok | :ignore
defdelegate accept_distributed_trace_headers(headers), to: NewRelic.DistributedTrace

@type name :: String.t() | {primary_name :: String.t(), secondary_name :: String.t()}

@doc """
Expand Down
67 changes: 44 additions & 23 deletions lib/new_relic/distributed_trace.ex
Original file line number Diff line number Diff line change
Expand Up @@ -12,7 +12,7 @@ defmodule NewRelic.DistributedTrace do
def start(type, headers \\ %{})

def start(:http, headers) do
if NewRelic.Config.feature?(:distributed_tracing) do
if NewRelic.Config.enabled?() && NewRelic.Config.feature?(:distributed_tracing) do
determine_context(headers)
|> track_transaction(transport_type: "HTTP")
end
Expand All @@ -21,22 +21,44 @@ defmodule NewRelic.DistributedTrace do
end

def start(:other, headers) do
if NewRelic.Config.feature?(:distributed_tracing) do
if NewRelic.Config.enabled?() && NewRelic.Config.feature?(:distributed_tracing) do
determine_context(headers)
|> track_transaction(transport_type: "Other")
end

:ok
end

def accept_distributed_trace_headers(headers) do
with true <- NewRelic.Config.enabled?() && NewRelic.Config.feature?(:distributed_tracing),
true <- Transaction.Sidecar.tracking?(),
false <- accepted_inbound_headers?(get_tracing_context()),
%Context{} = context <- extract_context(normalize_headers(headers)) do
track_transaction(context, transport_type: "Other")
:ok
else
_ -> :ignore
end
end

defp accepted_inbound_headers?(%Context{source: source}) when source != :new, do: true
defp accepted_inbound_headers?(_), do: false

defp normalize_headers(headers) when is_map(headers), do: headers

defp normalize_headers(headers) when is_list(headers),
do: Map.new(headers, fn {key, value} -> {to_string(key), value} end)

defp normalize_headers(_), do: %{}

defp determine_context(headers) do
case accept_distributed_trace_headers(headers) do
case extract_context(headers) do
%Context{} = context -> context
_ -> generate_new_context()
end
end

defp accept_distributed_trace_headers(headers) do
defp extract_context(headers) do
w3c_headers(headers) || newrelic_header(headers) || :no_payload
end

Expand All @@ -63,25 +85,24 @@ defmodule NewRelic.DistributedTrace do
end

def distributed_trace_headers(:http) do
case get_tracing_context() do
nil ->
[]

context ->
context = %{
context
| span_guid: get_current_span_guid(),
timestamp: System.system_time(:millisecond)
}

nr_header = NewRelic.DistributedTrace.NewRelicContext.generate(context)
{traceparent, tracestate} = NewRelic.DistributedTrace.W3CTraceContext.generate(context)

[
{@nr_header, nr_header},
{@w3c_traceparent, traceparent},
{@w3c_tracestate, tracestate}
]
with true <- NewRelic.Config.enabled?(),
%Context{} = context <- get_tracing_context() do
context = %{
context
| span_guid: get_current_span_guid(),
timestamp: System.system_time(:millisecond)
}

nr_header = NewRelic.DistributedTrace.NewRelicContext.generate(context)
{traceparent, tracestate} = NewRelic.DistributedTrace.W3CTraceContext.generate(context)

[
{@nr_header, nr_header},
{@w3c_traceparent, traceparent},
{@w3c_tracestate, tracestate}
]
else
_ -> []
end
end

Expand Down
16 changes: 12 additions & 4 deletions lib/new_relic/transaction/sidecar.ex
Original file line number Diff line number Diff line change
Expand Up @@ -83,13 +83,21 @@ defmodule NewRelic.Transaction.Sidecar do
end

def trace_context(context) do
:ets.insert(__MODULE__.ContextStore, {{:context, get_sidecar()}, context})
case get_sidecar() do
sidecar when is_pid(sidecar) ->
:ets.insert(__MODULE__.ContextStore, {{:context, sidecar}, context})

_ ->
:no_sidecar
end
end

def trace_context() do
case :ets.lookup(__MODULE__.ContextStore, {:context, get_sidecar()}) do
[{_, value}] -> value
[] -> nil
with sidecar when is_pid(sidecar) <- get_sidecar(),
[{_, value}] <- :ets.lookup(__MODULE__.ContextStore, {:context, sidecar}) do
value
else
_ -> nil
end
end

Expand Down
Loading
Loading