diff --git a/lib/anubis/server/session/scheduler.ex b/lib/anubis/server/session/scheduler.ex index 84225710..ca8a7894 100644 --- a/lib/anubis/server/session/scheduler.ex +++ b/lib/anubis/server/session/scheduler.ex @@ -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(), diff --git a/lib/anubis/server/session/tasks.ex b/lib/anubis/server/session/tasks.ex index f507b7ad..dc89372e 100644 --- a/lib/anubis/server/session/tasks.ex +++ b/lib/anubis/server/session/tasks.ex @@ -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) @@ -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) diff --git a/lib/anubis/telemetry.ex b/lib/anubis/telemetry.ex index 4788f667..d55af558 100644 --- a/lib/anubis/telemetry.ex +++ b/lib/anubis/telemetry.ex @@ -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. @@ -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 + 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 diff --git a/pages/testing.md b/pages/testing.md index 8ba72ea5..241f9e91 100644 --- a/pages/testing.md +++ b/pages/testing.md @@ -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)) + end, + nil +) +``` + +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. diff --git a/test/anubis/server/session_test.exs b/test/anubis/server/session_test.exs index fbd2fa12..3f0e76bc 100644 --- a/test/anubis/server/session_test.exs +++ b/test/anubis/server/session_test.exs @@ -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} + 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) diff --git a/test/anubis/server/tasks_test.exs b/test/anubis/server/tasks_test.exs index b4f0d93d..ab3f9979 100644 --- a/test/anubis/server/tasks_test.exs +++ b/test/anubis/server/tasks_test.exs @@ -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