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
16 changes: 4 additions & 12 deletions lib/anubis/server/session/scheduler.ex
Original file line number Diff line number Diff line change
Expand Up @@ -183,25 +183,17 @@ defmodule Anubis.Server.Session.Scheduler do

defp do_handle_request(module, %{"method" => "tools/call"} = request, frame, _method) do
tool_name = get_in(request, ["params", "name"])
arguments = get_in(request, ["params", "arguments"])

:telemetry.span(
[:anubis_mcp | Telemetry.event_server_tool_call()],
%{tool: tool_name},
fn ->
result = module.handle_request(request, frame)
{result, %{tool: tool_name, is_error: tool_call_error?(result)}}
end
)
Telemetry.span_tool_call(tool_name, arguments, fn ->
module.handle_request(request, frame)
end)
end

defp do_handle_request(module, request, frame, _method) do
module.handle_request(request, frame)
end

defp tool_call_error?({:error, _reason, _frame}), do: true
defp tool_call_error?({:reply, %{"isError" => true}, _frame}), do: true
defp tool_call_error?(_result), do: false

defp decode_task_result({:reply, response, %Frame{} = frame}, inflight, state) do
Telemetry.execute(
Telemetry.event_server_response(),
Expand Down
7 changes: 6 additions & 1 deletion lib/anubis/server/session/tasks.ex
Original file line number Diff line number Diff line change
Expand Up @@ -11,6 +11,7 @@ defmodule Anubis.Server.Session.Tasks do
alias Anubis.Server.Handlers.Tasks, as: TasksHandler
alias Anubis.Server.Session.ServerRequests
alias Anubis.Server.Task, as: McpTask
alias Anubis.Telemetry

@default_task_ttl 60_000
@max_task_ttl to_timeout(hour: 1)
Expand Down Expand Up @@ -401,10 +402,14 @@ defmodule Anubis.Server.Session.Tasks do
request = %{request | "params" => Map.delete(params, "task")}

server_module = state.server_module
tool_name = params["name"]
arguments = params["arguments"]

worker =
Task.Supervisor.async_nolink(state.task_supervisor, fn ->
Handlers.handle(request, server_module, frame)
Telemetry.span_tool_call(tool_name, arguments, fn ->
Handlers.handle(request, server_module, frame)
end)
end)

ttl_timer = Process.send_after(self(), {:task_expired, task.id}, ttl)
Expand Down
57 changes: 57 additions & 0 deletions lib/anubis/telemetry.ex
Original file line number Diff line number Diff line change
@@ -1,6 +1,8 @@
defmodule Anubis.Telemetry do
@moduledoc false

@default_capture_tool_payload false

@doc """
Execute a telemetry event with the Anubis MCP namespace.

Expand All @@ -14,6 +16,61 @@ defmodule Anubis.Telemetry do
:telemetry.execute([:anubis_mcp | event_name], measurements, metadata)
end

@doc """
Wraps a tool call handler invocation in the `[:server, :tool_call]`
telemetry span.

Shared by the synchronous scheduler dispatch and the task-augmented
`tools/call` worker path so both surface identical span data regardless
of which route a given request took.

The span's `:start` metadata always carries `tool`; the `:stop` metadata
always carries `tool` and `is_error`. When
`:telemetry_capture_tool_payload` is enabled (defaults to `false`), the
`:start` metadata also carries `arguments` and the `:stop` metadata also
carries `result`. See `pages/testing.md` for the rationale behind the
opt-in default.

## Examples

iex> Anubis.Telemetry.span_tool_call("get_weather", %{"city" => "NYC"}, fn -> :ok end)
:ok
"""
@spec span_tool_call(String.t() | nil, map() | nil, (-> result)) :: result when result: var
Comment thread
coderabbitai[bot] marked this conversation as resolved.
def span_tool_call(tool_name, arguments, fun) do
capture_payload? =
Application.get_env(:anubis_mcp, :telemetry_capture_tool_payload, @default_capture_tool_payload)

start_metadata =
if capture_payload? do
%{tool: tool_name, arguments: arguments}
else
%{tool: tool_name}
end

:telemetry.span(
[:anubis_mcp | event_server_tool_call()],
start_metadata,
fn ->
result = fun.()
is_error = tool_call_error?(result)

stop_metadata =
if capture_payload? do
%{tool: tool_name, is_error: is_error, result: result}
else
%{tool: tool_name, is_error: is_error}
end

{result, stop_metadata}
end
)
end

defp tool_call_error?({:error, _reason, _frame}), do: true
defp tool_call_error?({:reply, %{"isError" => true}, _frame}), do: true
defp tool_call_error?(_result), do: false

# Define event name constants to ensure consistency

# Client events
Expand Down
50 changes: 50 additions & 0 deletions pages/testing.md
Original file line number Diff line number Diff line change
Expand Up @@ -176,6 +176,56 @@ end

The `start: true` option forces the HTTP transport to boot even though no Phoenix endpoint is serving during tests. One integration test per server is usually enough; the per-component behavior belongs in the unit tests above.

## Observing tool calls

Every `tools/call` request — including task-augmented calls dispatched via
the `tasks/` worker path — is wrapped in a `:telemetry.span/3` under
`[:anubis_mcp, :server, :tool_call]`. By default the span's metadata carries
only the tool name and whether the call errored — enough to build dashboards
and alerts, not enough to answer "what was this client actually asking, and
what did we return."

Set `:telemetry_capture_tool_payload` to opt into the full payload:

```elixir
# config/config.exs (or runtime.exs)
config :anubis_mcp, :telemetry_capture_tool_payload, true
```

With the flag enabled, the `:start` event's metadata gains `arguments` (the
raw `params.arguments` map from the request) and the `:stop` event's
metadata gains `result` (the value your callback returned). `arguments` is
only present on `:start`, and `result`/`is_error` only on `:stop` — attach to
both events if you need both:

```elixir
require Logger

:telemetry.attach_many(
"log-tool-payloads",
[
[:anubis_mcp, :server, :tool_call, :start],
[:anubis_mcp, :server, :tool_call, :stop]
],
fn
[:anubis_mcp, :server, :tool_call, :start], _measurements, %{tool: tool, arguments: args}, _config ->
Logger.info("tool_call_start", tool: tool, arguments: args)

[:anubis_mcp, :server, :tool_call, :stop], _measurements, %{tool: tool, result: result}, _config ->
Logger.info("tool_call_stop", tool: tool, result: inspect(result))
Comment thread
coderabbitai[bot] marked this conversation as resolved.
end,
nil
)
Comment thread
coderabbitai[bot] marked this conversation as resolved.
```

The flag defaults to `false`. Tool arguments and results are the
highest-cardinality, highest-PII surface in the request lifecycle — this
mirrors the `opt_in` requirement level OpenTelemetry's GenAI semantic
conventions assign to the equivalent `gen_ai.tool.call.arguments` /
`gen_ai.tool.call.result` span attributes. Only enable it where you also
control retention and redaction of whatever sink receives the telemetry
event (a Datadog trace, a log pipeline, a database).

## Next steps

- [Building a Server](building-a-server.md) documents the callbacks tested here.
Expand Down
83 changes: 83 additions & 0 deletions test/anubis/server/session_test.exs
Original file line number Diff line number Diff line change
Expand Up @@ -744,6 +744,89 @@ defmodule Anubis.Server.SessionTest do
end
end

describe "tool_call telemetry opt-in payload capture" do
setup do
original_config = Application.fetch_env(:anubis_mcp, :telemetry_capture_tool_payload)

on_exit(fn ->
case original_config do
{:ok, value} -> Application.put_env(:anubis_mcp, :telemetry_capture_tool_payload, value)
:error -> Application.delete_env(:anubis_mcp, :telemetry_capture_tool_payload)
end
end)

test_pid = self()
handler_id = "test-tool-call-payload-#{System.unique_integer([:positive])}"

:telemetry.attach_many(
handler_id,
[
[:anubis_mcp, :server, :tool_call, :start],
[:anubis_mcp, :server, :tool_call, :stop]
],
fn
[:anubis_mcp, :server, :tool_call, :start], _m, meta, _c -> send(test_pid, {:tool_call_start, meta})
[:anubis_mcp, :server, :tool_call, :stop], _m, meta, _c -> send(test_pid, {:tool_call_stop, meta})
end,
nil
)

on_exit(fn -> :telemetry.detach(handler_id) end)

transport_name = Registry.transport_name(TasksStubServer, StubTransport)
start_supervised!({StubTransport, name: transport_name}, id: :payload_transport)
task_sup = Registry.task_supervisor_name(TasksStubServer)
start_supervised!({Task.Supervisor, name: task_sup}, id: :payload_task_sup)

session_id = "tool-call-payload-#{System.unique_integer([:positive])}"
session_name = Registry.session_name(TasksStubServer, session_id)

session =
start_supervised!(
{Session,
session_id: session_id,
server_module: TasksStubServer,
name: session_name,
transport: [layer: StubTransport, name: transport_name],
task_supervisor: task_sup},
id: :payload_session
)

init_msg = init_request("2025-03-26", %{"name" => "TestClient", "version" => "1.0.0"})
assert {:ok, _} = GenServer.call(session, {:mcp_request, init_msg, %{}})

init_notification = build_notification("notifications/initialized", %{})
assert :ok = GenServer.cast(session, {:mcp_notification, init_notification, %{}})

%{session: session}
Comment thread
acollado-g2 marked this conversation as resolved.
end

test "arguments/result are absent from span metadata when the flag is unset (default)", %{session: session} do
Application.delete_env(:anubis_mcp, :telemetry_capture_tool_payload)

request = build_request("tools/call", %{"name" => "no_tasks", "arguments" => %{"x" => 1}}, 1)

assert {:ok, _} = GenServer.call(session, {:mcp_request, request, %{}})

assert_receive {:tool_call_start, start_meta}, 500
refute Map.has_key?(start_meta, :arguments)

assert_receive {:tool_call_stop, stop_meta}, 500
refute Map.has_key?(stop_meta, :result)
end

test "arguments/result are present in span metadata when the flag is enabled", %{session: session} do
Application.put_env(:anubis_mcp, :telemetry_capture_tool_payload, true)

request = build_request("tools/call", %{"name" => "no_tasks", "arguments" => %{"x" => 1}}, 2)

assert {:ok, _} = GenServer.call(session, {:mcp_request, request, %{}})

assert_receive {:tool_call_start, %{tool: "no_tasks", arguments: %{"x" => 1}}}, 500
assert_receive {:tool_call_stop, %{tool: "no_tasks", is_error: false, result: _result}}, 500
end
end

describe "terminate/2 on supervisor-initiated stop" do
defp start_supervised_session(server_module, session_id) do
transport_name = Registry.transport_name(server_module, StubTransport)
Expand Down
30 changes: 30 additions & 0 deletions test/anubis/server/tasks_test.exs
Original file line number Diff line number Diff line change
Expand Up @@ -46,6 +46,36 @@ defmodule Anubis.Server.TasksTest do

assert decoded["error"]["code"] == -32_601
end

test "emits [:anubis_mcp, :server, :tool_call] telemetry for the worker's async execution", %{session: session} do
test_pid = self()
handler_id = "test-task-augmented-tool-call-#{System.unique_integer([:positive])}"

:telemetry.attach_many(
handler_id,
[
[:anubis_mcp, :server, :tool_call, :start],
[:anubis_mcp, :server, :tool_call, :stop]
],
fn
[:anubis_mcp, :server, :tool_call, :start], _m, meta, _c -> send(test_pid, {:task_tool_call_start, meta})
[:anubis_mcp, :server, :tool_call, :stop], _m, meta, _c -> send(test_pid, {:task_tool_call_stop, meta})
end,
nil
)

on_exit(fn -> :telemetry.detach(handler_id) end)

decoded = create_task_call(session, "always_fails", %{"reason" => "kaboom"}, "req-1")
task_id = decoded["result"]["task"]["taskId"]

# The worker runs asynchronously — wait for span :stop rather than
# asserting immediately after the CreateTaskResult reply.
assert_receive {:task_tool_call_start, %{tool: "always_fails"}}, 500
assert_receive {:task_tool_call_stop, %{tool: "always_fails", is_error: true}}, 500

SyncHelpers.await_state(session, fn state -> not Map.has_key?(state.tasks, task_id) end)
end
end

describe "tasks/get" do
Expand Down
Loading