diff --git a/lib/anubis/client/base.ex b/lib/anubis/client/base.ex index 23bd9869..d54d6cfd 100644 --- a/lib/anubis/client/base.ex +++ b/lib/anubis/client/base.ex @@ -112,6 +112,8 @@ defmodule Anubis.Client.Base do optional(:sampling | String.t()) => %{} } + @default_operation_timeout to_timeout(second: 30) + @typedoc """ MCP client initialization options @@ -136,7 +138,8 @@ defmodule Anubis.Client.Base do {:transport, {:required, {:custom, &Anubis.client_transport/1}}}, {:client_info, {:required, :map}}, {:capabilities, {:required, :map}}, - {:protocol_version, {:string, {:default, @default_protocol_version}}} + {:protocol_version, {:string, {:default, @default_protocol_version}}}, + {:timeout, {:integer, {:default, @default_operation_timeout}}} ]) @doc """ @@ -172,7 +175,7 @@ defmodule Anubis.Client.Base do method: "ping", params: %{}, progress_opts: Keyword.get(opts, :progress), - timeout: Keyword.get(opts, :timeout) + timeout: Keyword.get(opts, :timeout, @default_operation_timeout) }) buffer_timeout = operation.timeout + to_timeout(second: 1) @@ -200,7 +203,7 @@ defmodule Anubis.Client.Base do method: "resources/list", params: params, progress_opts: Keyword.get(opts, :progress), - timeout: Keyword.get(opts, :timeout) + timeout: Keyword.get(opts, :timeout, @default_operation_timeout) }) buffer_timeout = operation.timeout + to_timeout(second: 1) @@ -228,7 +231,7 @@ defmodule Anubis.Client.Base do method: "resources/templates/list", params: params, progress_opts: Keyword.get(opts, :progress), - timeout: Keyword.get(opts, :timeout) + timeout: Keyword.get(opts, :timeout, @default_operation_timeout) }) buffer_timeout = operation.timeout + to_timeout(second: 1) @@ -253,7 +256,7 @@ defmodule Anubis.Client.Base do method: "resources/read", params: %{"uri" => uri}, progress_opts: Keyword.get(opts, :progress), - timeout: Keyword.get(opts, :timeout) + timeout: Keyword.get(opts, :timeout, @default_operation_timeout) }) buffer_timeout = operation.timeout + to_timeout(second: 1) @@ -281,7 +284,7 @@ defmodule Anubis.Client.Base do method: "prompts/list", params: params, progress_opts: Keyword.get(opts, :progress), - timeout: Keyword.get(opts, :timeout) + timeout: Keyword.get(opts, :timeout, @default_operation_timeout) }) buffer_timeout = operation.timeout + to_timeout(second: 1) @@ -309,7 +312,7 @@ defmodule Anubis.Client.Base do method: "prompts/get", params: params, progress_opts: Keyword.get(opts, :progress), - timeout: Keyword.get(opts, :timeout) + timeout: Keyword.get(opts, :timeout, @default_operation_timeout) }) buffer_timeout = operation.timeout + to_timeout(second: 1) @@ -337,7 +340,7 @@ defmodule Anubis.Client.Base do method: "tools/list", params: params, progress_opts: Keyword.get(opts, :progress), - timeout: Keyword.get(opts, :timeout) + timeout: Keyword.get(opts, :timeout, @default_operation_timeout) }) buffer_timeout = operation.timeout + to_timeout(second: 1) @@ -365,7 +368,7 @@ defmodule Anubis.Client.Base do method: "tools/call", params: params, progress_opts: Keyword.get(opts, :progress), - timeout: Keyword.get(opts, :timeout) + timeout: Keyword.get(opts, :timeout, @default_operation_timeout) }) buffer_timeout = operation.timeout + to_timeout(second: 1) @@ -418,7 +421,8 @@ defmodule Anubis.Client.Base do operation = Operation.new(%{ method: "logging/setLevel", - params: %{"level" => level} + params: %{"level" => level}, + timeout: @default_operation_timeout }) buffer_timeout = operation.timeout + to_timeout(second: 1) @@ -474,7 +478,7 @@ defmodule Anubis.Client.Base do method: "completion/complete", params: params, progress_opts: Keyword.get(opts, :progress), - timeout: Keyword.get(opts, :timeout) + timeout: Keyword.get(opts, :timeout, @default_operation_timeout) }) buffer_timeout = operation.timeout + to_timeout(second: 1) @@ -785,7 +789,8 @@ defmodule Anubis.Client.Base do client_info: opts.client_info, capabilities: opts.capabilities, protocol_version: protocol_version, - transport: transport + transport: transport, + timeout: opts.timeout }) client_name = get_in(opts, [:client_info, "name"]) @@ -827,7 +832,7 @@ defmodule Anubis.Client.Base do {request_id, updated_state} = State.add_request_from_operation(state, operation, from), {:ok, request_data} <- encode_request(method, params_with_token, request_id), - :ok <- send_to_transport(state.transport, request_data) do + :ok <- send_to_transport(state.transport, request_data, timeout: operation.timeout) do Telemetry.execute( Telemetry.event_client_request(), %{system_time: System.system_time()}, @@ -885,7 +890,7 @@ defmodule Anubis.Client.Base do "progress" => progress, "total" => total }) do - send_to_transport(state.transport, notification) + send_to_transport(state.transport, notification, timeout: state.timeout) end, state} end @@ -972,14 +977,15 @@ defmodule Anubis.Client.Base do operation = Operation.new(%{ method: "initialize", - params: params + params: params, + timeout: state.timeout }) {request_id, updated_state} = State.add_request_from_operation(state, operation, {self(), make_ref()}) with {:ok, request_data} <- encode_request("initialize", params, request_id), - :ok <- send_to_transport(state.transport, request_data) do + :ok <- send_to_transport(state.transport, request_data, timeout: operation.timeout) do {:noreply, updated_state} else err -> {:stop, err, state} @@ -1018,7 +1024,7 @@ defmodule Anubis.Client.Base do with {:ok, response_data} <- Message.encode_response(%{"result" => roots_result}, id), - :ok <- send_to_transport(state.transport, response_data) do + :ok <- send_to_transport(state.transport, response_data, timeout: state.timeout) do Logging.client_event("roots_list_request", %{id: id, roots_count: roots_count}) Telemetry.execute( @@ -1044,7 +1050,7 @@ defmodule Anubis.Client.Base do defp handle_server_request(%{"method" => "ping", "id" => id}, state) do with {:ok, response_data} <- Message.encode_response(%{"result" => %{}}, id), - :ok <- send_to_transport(state.transport, response_data) do + :ok <- send_to_transport(state.transport, response_data, timeout: state.timeout) do {:noreply, state} else err -> @@ -1494,15 +1500,15 @@ defmodule Anubis.Client.Base do send_notification(state, "notifications/cancelled", params) end - defp send_to_transport(transport, data) do - with {:error, reason} <- transport.layer.send_message(transport.name, data) do + defp send_to_transport(transport, data, opts) do + with {:error, reason} <- transport.layer.send_message(transport.name, data, opts) do {:error, Error.transport(:send_failure, %{original_reason: reason})} end end defp send_notification(state, method, params \\ %{}) do with {:ok, notification_data} <- encode_notification(method, params) do - send_to_transport(state.transport, notification_data) + send_to_transport(state.transport, notification_data, timeout: state.timeout) end end @@ -1593,7 +1599,7 @@ defmodule Anubis.Client.Base do defp send_sampling_response(id, response, state) do transport = state.transport - :ok = transport.layer.send_message(transport.name, response) + :ok = transport.layer.send_message(transport.name, response, timeout: state.timeout) Telemetry.execute( Telemetry.event_client_response(), @@ -1605,7 +1611,7 @@ defmodule Anubis.Client.Base do defp send_sampling_error(id, message, code, reason, %{transport: transport} = state) do error = %Error{code: -1, message: message, data: %{"reason" => reason}} {:ok, response} = Error.to_json_rpc(error, id) - :ok = transport.layer.send_message(transport.name, response) + :ok = transport.layer.send_message(transport.name, response, timeout: state.timeout) Logging.client_event( "sampling_error", diff --git a/lib/anubis/client/operation.ex b/lib/anubis/client/operation.ex index 6c7041b7..b0a7fb0f 100644 --- a/lib/anubis/client/operation.ex +++ b/lib/anubis/client/operation.ex @@ -9,8 +9,6 @@ defmodule Anubis.Client.Operation do - `timeout` - The timeout for this specific operation (default: 30 seconds) """ - @default_timeout to_timeout(second: 30) - @type progress_options :: [ token: String.t() | integer(), callback: (String.t() | integer(), number(), number() | nil -> any()) @@ -25,9 +23,9 @@ defmodule Anubis.Client.Operation do defstruct [ :method, + :timeout, params: %{}, - progress_opts: [], - timeout: @default_timeout + progress_opts: [] ] @doc """ @@ -47,12 +45,12 @@ defmodule Anubis.Client.Operation do optional(:progress_opts) => progress_options() | nil, optional(:timeout) => pos_integer() }) :: t() - def new(%{method: method} = attrs) do + def new(%{method: method, timeout: timeout} = attrs) do %__MODULE__{ method: method, params: Map.get(attrs, :params) || %{}, progress_opts: Map.get(attrs, :progress_opts), - timeout: Map.get(attrs, :timeout) || @default_timeout + timeout: timeout } end end diff --git a/lib/anubis/client/state.ex b/lib/anubis/client/state.ex index 0d93fff2..52338a12 100644 --- a/lib/anubis/client/state.ex +++ b/lib/anubis/client/state.ex @@ -14,6 +14,7 @@ defmodule Anubis.Client.State do server_capabilities: map() | nil, server_info: map() | nil, protocol_version: String.t(), + timeout: pos_integer(), transport: map(), pending_requests: %{String.t() => Request.t()}, progress_callbacks: %{String.t() => Base.progress_callback()}, @@ -28,6 +29,7 @@ defmodule Anubis.Client.State do :capabilities, :server_capabilities, :server_info, + :timeout, :protocol_version, :transport, pending_requests: %{}, @@ -43,7 +45,8 @@ defmodule Anubis.Client.State do client_info: opts.client_info, capabilities: opts.capabilities, protocol_version: opts.protocol_version, - transport: opts.transport + transport: opts.transport, + timeout: opts.timeout } end diff --git a/lib/anubis/server/base.ex b/lib/anubis/server/base.ex index 419c92c8..ebd244e0 100644 --- a/lib/anubis/server/base.ex +++ b/lib/anubis/server/base.ex @@ -60,7 +60,8 @@ defmodule Anubis.Server.Base do {:name, {:required, {:custom, &Anubis.genserver_name/1}}}, {:transport, {:required, {:custom, &Anubis.server_transport/1}}}, {:registry, {:atom, {:default, Anubis.Server.Registry}}}, - {:session_idle_timeout, {{:integer, {:gte, 1}}, {:default, @default_session_idle_timeout}}} + {:session_idle_timeout, {{:integer, {:gte, 1}}, {:default, @default_session_idle_timeout}}}, + {:timeout, {:integer, {:default, to_timeout(second: 30)}}} ]) @spec start_link(Enumerable.t(option())) :: GenServer.on_start() @@ -90,7 +91,8 @@ defmodule Anubis.Server.Base do session_idle_timeout: opts.session_idle_timeout, expiry_timers: %{}, frame: Frame.new(), - server_requests: %{} + server_requests: %{}, + timeout: opts.timeout } Logging.server_event("starting", %{ @@ -233,7 +235,7 @@ defmodule Anubis.Server.Base do def handle_info({:send_notification, method, params}, state) do with {:ok, notification} <- encode_notification(method, params), - :ok <- send_to_transport(state.transport, notification) do + :ok <- send_to_transport(state.transport, notification, timeout: state.timeout) do {:noreply, state} else {:error, err} -> @@ -674,12 +676,12 @@ defmodule Anubis.Server.Base do Message.encode_notification(notification) end - defp send_to_transport(nil, _data) do + defp send_to_transport(nil, _data, _opts) do {:error, Error.transport(:no_transport, %{message: "No transport configured"})} end - defp send_to_transport(%{layer: layer, name: name}, data) do - with {:error, reason} <- layer.send_message(name, data) do + defp send_to_transport(%{layer: layer, name: name}, data, opts) do + with {:error, reason} <- layer.send_message(name, data, opts) do {:error, Error.transport(:send_failure, %{original_reason: reason})} end end @@ -723,7 +725,7 @@ defmodule Anubis.Server.Base do with :ok <- validate_client_capability(state, "sampling"), {:ok, request_data} <- encode_request("sampling/createMessage", params, request_id), - :ok <- send_to_transport(state.transport, request_data) do + :ok <- send_to_transport(state.transport, request_data, timeout: state.timeout) do Logging.server_event("sent_sampling_request", %{request_id: request_id}) {:noreply, state} else @@ -854,7 +856,7 @@ defmodule Anubis.Server.Base do with :ok <- validate_client_capability(state, "roots"), {:ok, request_data} <- encode_request("roots/list", %{}, request_id), - :ok <- send_to_transport(state.transport, request_data) do + :ok <- send_to_transport(state.transport, request_data, timeout: state.timeout) do Logging.server_event("sent_roots_request", %{request_id: request_id}) {:noreply, state} else @@ -890,7 +892,7 @@ defmodule Anubis.Server.Base do "requestId" => request_id, "reason" => "timeout" }), - :ok <- send_to_transport(state.transport, notification) do + :ok <- send_to_transport(state.transport, notification, timeout: state.timeout) do Logging.server_event( "roots_request_timeout_cancelled", %{request_id: request_id} diff --git a/lib/anubis/server/transport/sse.ex b/lib/anubis/server/transport/sse.ex index d33a156d..cfc1a4e9 100644 --- a/lib/anubis/server/transport/sse.ex +++ b/lib/anubis/server/transport/sse.ex @@ -131,9 +131,8 @@ defmodule Anubis.Server.Transport.SSE do * `{:error, reason}` otherwise """ @impl Transport - @spec send_message(GenServer.server(), binary()) :: :ok | {:error, term()} - def send_message(transport, message) when is_binary(message) do - GenServer.call(transport, {:send_message, message}) + def send_message(transport, message, opts) when is_binary(message) do + GenServer.call(transport, {:send_message, message}, opts[:timeout]) end @doc """ diff --git a/lib/anubis/server/transport/stdio.ex b/lib/anubis/server/transport/stdio.ex index 2802a99d..4847ef42 100644 --- a/lib/anubis/server/transport/stdio.ex +++ b/lib/anubis/server/transport/stdio.ex @@ -77,9 +77,8 @@ defmodule Anubis.Server.Transport.STDIO do * `{:error, reason}` otherwise """ @impl Transport - @spec send_message(GenServer.server(), binary()) :: :ok | {:error, term()} - def send_message(transport, message) when is_binary(message) do - GenServer.cast(transport, {:send, message}) + def send_message(transport, message, opts) when is_binary(message) do + GenServer.call(transport, {:send, message}, opts[:timeout]) end @doc """ @@ -262,7 +261,9 @@ defmodule Anubis.Server.Transport.STDIO do else case GenServer.call(server, {:request, message, "stdio", context}, timeout) do {:ok, response} when is_binary(response) -> - send_message(self(), response) + # send_message(self(), response) + # NOTE: will be fixed soon, we need to rewrite stdio for server + :ok {:error, reason} -> Logging.transport_event("server_error", %{reason: reason}, level: :error) diff --git a/lib/anubis/server/transport/streamable_http.ex b/lib/anubis/server/transport/streamable_http.ex index da285310..01f1f273 100644 --- a/lib/anubis/server/transport/streamable_http.ex +++ b/lib/anubis/server/transport/streamable_http.ex @@ -118,9 +118,8 @@ defmodule Anubis.Server.Transport.StreamableHTTP do * `{:error, reason}` otherwise """ @impl Transport - @spec send_message(GenServer.server(), binary()) :: :ok | {:error, term()} - def send_message(transport, message) when is_binary(message) do - GenServer.call(transport, {:send_message, message}, 5000) + def send_message(transport, message, opts) when is_binary(message) do + GenServer.call(transport, {:send_message, message}, opts[:timeout]) end @doc """ diff --git a/lib/anubis/transport/behaviour.ex b/lib/anubis/transport/behaviour.ex index bdd2bf55..3267f36b 100644 --- a/lib/anubis/transport/behaviour.ex +++ b/lib/anubis/transport/behaviour.ex @@ -11,12 +11,10 @@ defmodule Anubis.Transport.Behaviour do @type reason :: term() | Error.t() @callback start_link(keyword()) :: GenServer.on_start() - @callback send_message(t(), message()) :: :ok | {:error, reason()} - @callback send_message(t(), message(), keyword()) :: :ok | {:error, reason()} + @callback send_message(t(), message(), list(opt)) :: :ok | {:error, reason()} + when opt: {:timeout, pos_integer()} @callback shutdown(t()) :: :ok | {:error, reason()} - @optional_callbacks send_message: 3 - @doc """ Returns the list of MCP protocol versions supported by this transport. diff --git a/lib/anubis/transport/sse.ex b/lib/anubis/transport/sse.ex index 4127497f..8704bea8 100644 --- a/lib/anubis/transport/sse.ex +++ b/lib/anubis/transport/sse.ex @@ -101,8 +101,8 @@ defmodule Anubis.Transport.SSE do end @impl Transport - def send_message(pid, message) when is_binary(message) do - GenServer.call(pid, {:send, message}) + def send_message(pid, message, opts) when is_binary(message) do + GenServer.call(pid, {:send, message}, opts[:timeout]) end @impl Transport diff --git a/lib/anubis/transport/stdio.ex b/lib/anubis/transport/stdio.ex index 52ad2a5f..88c55a7f 100644 --- a/lib/anubis/transport/stdio.ex +++ b/lib/anubis/transport/stdio.ex @@ -73,8 +73,8 @@ defmodule Anubis.Transport.STDIO do end @impl Transport - def send_message(pid \\ __MODULE__, message) when is_binary(message) do - GenServer.call(pid, {:send, message}) + def send_message(pid \\ __MODULE__, message, opts) when is_binary(message) do + GenServer.call(pid, {:send, message}, opts[:timeout]) end @impl Transport diff --git a/lib/anubis/transport/streamable_http.ex b/lib/anubis/transport/streamable_http.ex index 28b73a98..e171c060 100644 --- a/lib/anubis/transport/streamable_http.ex +++ b/lib/anubis/transport/streamable_http.ex @@ -96,8 +96,8 @@ defmodule Anubis.Transport.StreamableHTTP do end @impl Transport - def send_message(pid \\ __MODULE__, message) when is_binary(message) do - GenServer.call(pid, {:send, message}) + def send_message(pid \\ __MODULE__, message, opts) when is_binary(message) do + GenServer.call(pid, {:send, message}, opts[:timeout]) end @impl Transport diff --git a/lib/anubis/transport/websocket.ex b/lib/anubis/transport/websocket.ex index e475b4cc..f8791e84 100644 --- a/lib/anubis/transport/websocket.ex +++ b/lib/anubis/transport/websocket.ex @@ -81,8 +81,8 @@ if Code.ensure_loaded?(:gun) do end @impl Transport - def send_message(pid, message) when is_binary(message) do - GenServer.call(pid, {:send, message}) + def send_message(pid, message, opts) when is_binary(message) do + GenServer.call(pid, {:send, message}, opts[:timeout]) end @impl Transport diff --git a/test/anubis/client/base_test.exs b/test/anubis/client/base_test.exs index d2480661..25d2f888 100644 --- a/test/anubis/client/base_test.exs +++ b/test/anubis/client/base_test.exs @@ -18,7 +18,7 @@ defmodule Anubis.Client.BaseTest do describe "start_link/1" do test "starts the client with proper initialization" do - expect(Anubis.MockTransport, :send_message, fn _, message -> + expect(Anubis.MockTransport, :send_message, fn _, message, _ -> assert String.contains?(message, "initialize") assert String.contains?(message, "protocolVersion") assert String.contains?(message, "capabilities") @@ -46,7 +46,7 @@ defmodule Anubis.Client.BaseTest do setup :initialized_client test "ping sends correct request", %{client: client} do - expect(Anubis.MockTransport, :send_message, fn _, message -> + expect(Anubis.MockTransport, :send_message, fn _, message, _ -> decoded = JSON.decode!(message) assert decoded["method"] == "ping" assert decoded["params"] == %{} @@ -69,7 +69,7 @@ defmodule Anubis.Client.BaseTest do end test "list_resources sends correct request", %{client: client} do - expect(Anubis.MockTransport, :send_message, fn _, message -> + expect(Anubis.MockTransport, :send_message, fn _, message, _ -> decoded = JSON.decode!(message) assert decoded["method"] == "resources/list" assert decoded["params"] == %{} @@ -99,7 +99,7 @@ defmodule Anubis.Client.BaseTest do end test "list_resource_templates sends correct request", %{client: client} do - expect(Anubis.MockTransport, :send_message, fn _, message -> + expect(Anubis.MockTransport, :send_message, fn _, message, _ -> decoded = JSON.decode!(message) assert decoded["method"] == "resources/templates/list" assert decoded["params"] == %{} @@ -129,7 +129,7 @@ defmodule Anubis.Client.BaseTest do end test "list_resources with cursor", %{client: client} do - expect(Anubis.MockTransport, :send_message, fn _, message -> + expect(Anubis.MockTransport, :send_message, fn _, message, _ -> decoded = JSON.decode!(message) assert decoded["method"] == "resources/list" assert decoded["params"] == %{"cursor" => "next-page"} @@ -162,7 +162,7 @@ defmodule Anubis.Client.BaseTest do end test "read_resource sends correct request", %{client: client} do - expect(Anubis.MockTransport, :send_message, fn _, message -> + expect(Anubis.MockTransport, :send_message, fn _, message, _ -> decoded = JSON.decode!(message) assert decoded["method"] == "resources/read" assert decoded["params"] == %{"uri" => "test://uri"} @@ -192,7 +192,7 @@ defmodule Anubis.Client.BaseTest do end test "list_prompts sends correct request", %{client: client} do - expect(Anubis.MockTransport, :send_message, fn _, message -> + expect(Anubis.MockTransport, :send_message, fn _, message, _ -> decoded = JSON.decode!(message) assert decoded["method"] == "prompts/list" assert decoded["params"] == %{} @@ -222,7 +222,7 @@ defmodule Anubis.Client.BaseTest do end test "get_prompt sends correct request", %{client: client} do - expect(Anubis.MockTransport, :send_message, fn _, message -> + expect(Anubis.MockTransport, :send_message, fn _, message, _ -> decoded = JSON.decode!(message) assert decoded["method"] == "prompts/get" @@ -262,7 +262,7 @@ defmodule Anubis.Client.BaseTest do end test "list_tools sends correct request", %{client: client} do - expect(Anubis.MockTransport, :send_message, fn _, message -> + expect(Anubis.MockTransport, :send_message, fn _, message, _ -> decoded = JSON.decode!(message) assert decoded["method"] == "tools/list" assert decoded["params"] == %{} @@ -292,7 +292,7 @@ defmodule Anubis.Client.BaseTest do end test "call_tool sends correct request", %{client: client} do - expect(Anubis.MockTransport, :send_message, fn _, message -> + expect(Anubis.MockTransport, :send_message, fn _, message, _ -> decoded = JSON.decode!(message) assert decoded["method"] == "tools/call" @@ -330,7 +330,7 @@ defmodule Anubis.Client.BaseTest do end test "handles domain error responses as {:ok, response}", %{client: client} do - expect(Anubis.MockTransport, :send_message, fn _, message -> + expect(Anubis.MockTransport, :send_message, fn _, message, _ -> decoded = JSON.decode!(message) assert decoded["method"] == "tools/call" @@ -376,7 +376,7 @@ defmodule Anubis.Client.BaseTest do setup :initialized_client test "ping sends correct request since it is always supported", %{client: client} do - expect(Anubis.MockTransport, :send_message, fn _, message -> + expect(Anubis.MockTransport, :send_message, fn _, message, _ -> decoded = JSON.decode!(message) assert decoded["method"] == "ping" assert decoded["params"] == %{} @@ -411,7 +411,7 @@ defmodule Anubis.Client.BaseTest do setup :initialized_client test "handles error response", %{client: client} do - expect(Anubis.MockTransport, :send_message, fn _, message -> + expect(Anubis.MockTransport, :send_message, fn _, message, _ -> decoded = JSON.decode!(message) assert decoded["method"] == "ping" :ok @@ -434,7 +434,7 @@ defmodule Anubis.Client.BaseTest do end test "handles transport error", %{client: client} do - expect(Anubis.MockTransport, :send_message, fn _, _message -> + expect(Anubis.MockTransport, :send_message, fn _, _, _ -> {:error, :connection_closed} end) @@ -446,7 +446,7 @@ defmodule Anubis.Client.BaseTest do describe "capability management" do test "merge_capabilities correctly merges capabilities" do - expect(Anubis.MockTransport, :send_message, fn _, _message -> :ok end) + expect(Anubis.MockTransport, :send_message, fn _, _message, _ -> :ok end) client = start_supervised!( @@ -542,7 +542,7 @@ defmodule Anubis.Client.BaseTest do test "request with progress token includes it in params", %{client: client} do progress_token = "request_token_test" - expect(Anubis.MockTransport, :send_message, fn _, message -> + expect(Anubis.MockTransport, :send_message, fn _, message, _ -> decoded = JSON.decode!(message) assert decoded["method"] == "resources/list" @@ -594,7 +594,7 @@ defmodule Anubis.Client.BaseTest do } test "set_log_level sends the correct request", %{client: client} do - expect(Anubis.MockTransport, :send_message, fn _, message -> + expect(Anubis.MockTransport, :send_message, fn _, message, _ -> decoded = JSON.decode!(message) assert decoded["method"] == "logging/setLevel" assert decoded["params"]["level"] == "info" @@ -626,7 +626,7 @@ defmodule Anubis.Client.BaseTest do ref = %{"type" => "ref/prompt", "name" => "code_review"} argument = %{"name" => "language", "value" => "py"} - expect(Anubis.MockTransport, :send_message, fn _, message -> + expect(Anubis.MockTransport, :send_message, fn _, message, _ -> decoded = JSON.decode!(message) assert decoded["method"] == "completion/complete" assert decoded["params"]["ref"]["type"] == "ref/prompt" @@ -669,7 +669,7 @@ defmodule Anubis.Client.BaseTest do ref = %{"type" => "ref/resource", "uri" => "file:///path/to/file.txt"} argument = %{"name" => "encoding", "value" => "ut"} - expect(Anubis.MockTransport, :send_message, fn _, message -> + expect(Anubis.MockTransport, :send_message, fn _, message, _ -> decoded = JSON.decode!(message) assert decoded["method"] == "completion/complete" assert decoded["params"]["ref"]["type"] == "ref/resource" @@ -738,14 +738,14 @@ defmodule Anubis.Client.BaseTest do test "sends initialized notification after init" do Anubis.MockTransport # the handle_continue - |> expect(:send_message, fn _, message -> + |> expect(:send_message, fn _, message, _ -> decoded = JSON.decode!(message) assert decoded["method"] == "initialize" assert decoded["jsonrpc"] == "2.0" :ok end) # the send_notification - |> expect(:send_message, fn _, message -> + |> expect(:send_message, fn _, message, _ -> decoded = JSON.decode!(message) assert decoded["method"] == "notifications/initialized" :ok @@ -775,7 +775,7 @@ defmodule Anubis.Client.BaseTest do setup :initialized_client test "handles cancelled notification from server", %{client: client} do - expect(Anubis.MockTransport, :send_message, fn _, message -> + expect(Anubis.MockTransport, :send_message, fn _, message, _ -> decoded = JSON.decode!(message) assert decoded["method"] == "tools/call" :ok @@ -803,7 +803,7 @@ defmodule Anubis.Client.BaseTest do end test "client can cancel a request", %{client: client} do - expect(Anubis.MockTransport, :send_message, fn _, message -> + expect(Anubis.MockTransport, :send_message, fn _, message, _ -> decoded = JSON.decode!(message) assert decoded["method"] == "resources/list" :ok @@ -816,7 +816,7 @@ defmodule Anubis.Client.BaseTest do request_id = get_request_id(client, "resources/list") assert request_id - expect(Anubis.MockTransport, :send_message, fn _, message -> + expect(Anubis.MockTransport, :send_message, fn _, message, _ -> decoded = JSON.decode!(message) assert decoded["method"] == "notifications/cancelled" assert decoded["params"]["requestId"] == request_id @@ -847,7 +847,7 @@ defmodule Anubis.Client.BaseTest do end test "cancel_all_requests cancels all pending requests", %{client: client} do - expect(Anubis.MockTransport, :send_message, 2, fn _, message -> + expect(Anubis.MockTransport, :send_message, 2, fn _, message, _ -> decoded = JSON.decode!(message) assert decoded["method"] in ["resources/list", "tools/list"] :ok @@ -862,7 +862,7 @@ defmodule Anubis.Client.BaseTest do pending_count = map_size(state.pending_requests) assert pending_count == 2 - expect(Anubis.MockTransport, :send_message, 2, fn _, message -> + expect(Anubis.MockTransport, :send_message, 2, fn _, message, _ -> decoded = JSON.decode!(message) assert decoded["method"] == "notifications/cancelled" assert decoded["params"]["reason"] == "batch cancellation" @@ -887,13 +887,13 @@ defmodule Anubis.Client.BaseTest do test "request timeout sends cancellation notification", %{client: client} do test_timeout = 50 - expect(Anubis.MockTransport, :send_message, fn _, message -> + expect(Anubis.MockTransport, :send_message, fn _, message, _ -> decoded = JSON.decode!(message) assert decoded["method"] == "resources/list" :ok end) - expect(Anubis.MockTransport, :send_message, fn _, message -> + expect(Anubis.MockTransport, :send_message, fn _, message, _ -> decoded = JSON.decode!(message) assert decoded["method"] == "notifications/cancelled" assert decoded["params"]["reason"] == "timeout" @@ -918,14 +918,14 @@ defmodule Anubis.Client.BaseTest do %{client: client} do test_timeout = 50 - expect(Anubis.MockTransport, :send_message, fn _, message -> + expect(Anubis.MockTransport, :send_message, fn _, message, _ -> decoded = JSON.decode!(message) assert decoded["method"] == "resources/list" Process.sleep(test_timeout + 10) :ok end) - expect(Anubis.MockTransport, :send_message, fn _, _ -> :ok end) + expect(Anubis.MockTransport, :send_message, fn _, _, _ -> :ok end) Process.flag(:trap_exit, true) @@ -942,13 +942,13 @@ defmodule Anubis.Client.BaseTest do end test "client.close sends cancellation for pending requests", %{client: client} do - expect(Anubis.MockTransport, :send_message, fn _, message -> + expect(Anubis.MockTransport, :send_message, fn _, message, _ -> decoded = JSON.decode!(message) assert decoded["method"] == "resources/list" :ok end) - expect(Anubis.MockTransport, :send_message, fn _, message -> + expect(Anubis.MockTransport, :send_message, fn _, message, _ -> decoded = JSON.decode!(message) assert decoded["method"] == "notifications/cancelled" assert decoded["params"]["reason"] == "client closed" @@ -1099,7 +1099,7 @@ defmodule Anubis.Client.BaseTest do request_id = "server_req_123" - expect(Anubis.MockTransport, :send_message, fn _, message -> + expect(Anubis.MockTransport, :send_message, fn _, message, _ -> decoded = JSON.decode!(message) assert decoded["jsonrpc"] == "2.0" assert decoded["id"] == request_id @@ -1181,7 +1181,7 @@ defmodule Anubis.Client.BaseTest do "modelPreferences" => %{} } - expect(Anubis.MockTransport, :send_message, fn _, message -> + expect(Anubis.MockTransport, :send_message, fn _, message, _ -> decoded = JSON.decode!(message) assert decoded["jsonrpc"] == "2.0" assert decoded["id"] == request_id @@ -1218,7 +1218,7 @@ defmodule Anubis.Client.BaseTest do "modelPreferences" => %{} } - expect(Anubis.MockTransport, :send_message, fn _, message -> + expect(Anubis.MockTransport, :send_message, fn _, message, _ -> decoded = JSON.decode!(message) assert decoded["jsonrpc"] == "2.0" assert decoded["id"] == request_id @@ -1256,7 +1256,7 @@ defmodule Anubis.Client.BaseTest do request_id = "server_sampling_req_789" params = %{"messages" => [], "modelPreferences" => %{}} - expect(Anubis.MockTransport, :send_message, fn _, message -> + expect(Anubis.MockTransport, :send_message, fn _, message, _ -> decoded = JSON.decode!(message) assert decoded["jsonrpc"] == "2.0" assert decoded["id"] == request_id @@ -1294,7 +1294,7 @@ defmodule Anubis.Client.BaseTest do request_id = "server_sampling_req_999" params = %{"messages" => [], "modelPreferences" => %{}} - expect(Anubis.MockTransport, :send_message, fn _, message -> + expect(Anubis.MockTransport, :send_message, fn _, message, _ -> decoded = JSON.decode!(message) assert decoded["jsonrpc"] == "2.0" assert decoded["id"] == request_id @@ -1325,7 +1325,7 @@ defmodule Anubis.Client.BaseTest do @tag client_capabilities: %{"roots" => %{"listChanged" => true}} test "sends notification when adding a root", %{client: client} do - expect(Anubis.MockTransport, :send_message, fn _, message -> + expect(Anubis.MockTransport, :send_message, fn _, message, _ -> decoded = JSON.decode!(message) assert decoded["jsonrpc"] == "2.0" assert decoded["method"] == "notifications/roots/list_changed" @@ -1348,7 +1348,7 @@ defmodule Anubis.Client.BaseTest do Process.sleep(50) - expect(Anubis.MockTransport, :send_message, fn _, message -> + expect(Anubis.MockTransport, :send_message, fn _, message, _ -> decoded = JSON.decode!(message) assert decoded["jsonrpc"] == "2.0" assert decoded["method"] == "notifications/roots/list_changed" @@ -1380,7 +1380,7 @@ defmodule Anubis.Client.BaseTest do Process.sleep(50) - expect(Anubis.MockTransport, :send_message, fn _, message -> + expect(Anubis.MockTransport, :send_message, fn _, message, _ -> decoded = JSON.decode!(message) assert decoded["jsonrpc"] == "2.0" assert decoded["method"] == "notifications/roots/list_changed" @@ -1409,7 +1409,7 @@ defmodule Anubis.Client.BaseTest do setup :initialized_client test "validates tool call output when structuredContent is present", %{client: client} do - expect(Anubis.MockTransport, :send_message, fn _, message -> + expect(Anubis.MockTransport, :send_message, fn _, message, _ -> decoded = JSON.decode!(message) assert decoded["method"] == "tools/list" :ok @@ -1438,7 +1438,7 @@ defmodule Anubis.Client.BaseTest do send_response(client, response) assert {:ok, _} = Task.await(list_task) - expect(Anubis.MockTransport, :send_message, fn _, message -> + expect(Anubis.MockTransport, :send_message, fn _, message, _ -> decoded = JSON.decode!(message) assert decoded["method"] == "tools/call" assert decoded["params"]["name"] == "get_weather" @@ -1480,7 +1480,7 @@ defmodule Anubis.Client.BaseTest do assert {:ok, response} = Task.await(call_task) assert response.result["structuredContent"] == valid_structured - expect(Anubis.MockTransport, :send_message, fn _, _ -> :ok end) + expect(Anubis.MockTransport, :send_message, fn _, _, _ -> :ok end) invalid_task = Task.async(fn -> @@ -1515,7 +1515,7 @@ defmodule Anubis.Client.BaseTest do end test "handles tools with complex outputSchema", %{client: client} do - expect(Anubis.MockTransport, :send_message, fn _, _ -> :ok end) + expect(Anubis.MockTransport, :send_message, fn _, _, _ -> :ok end) task = Task.async(fn -> Anubis.Client.Base.list_tools(client) end) Process.sleep(50) diff --git a/test/anubis/client/state_test.exs b/test/anubis/client/state_test.exs index 0036605c..f964327d 100644 --- a/test/anubis/client/state_test.exs +++ b/test/anubis/client/state_test.exs @@ -12,7 +12,8 @@ defmodule Anubis.Client.StateTest do client_info: %{"name" => "TestClient", "version" => "1.0.0"}, capabilities: %{"resources" => %{}}, protocol_version: "2024-11-05", - transport: %{layer: :fake_transport, name: :fake_name} + transport: %{layer: :fake_transport, name: :fake_name}, + timeout: 30_000 } state = State.new(opts) @@ -21,6 +22,7 @@ defmodule Anubis.Client.StateTest do assert state.capabilities == %{"resources" => %{}} assert state.protocol_version == "2024-11-05" assert state.transport == %{layer: :fake_transport, name: :fake_name} + assert state.timeout == 30_000 assert state.pending_requests == %{} assert state.progress_callbacks == %{} assert state.log_callback == nil @@ -34,7 +36,8 @@ defmodule Anubis.Client.StateTest do operation = Operation.new(%{ - method: "test_method" + method: "test_method", + timeout: 30_000 }) {request_id, updated_state} = @@ -58,7 +61,7 @@ defmodule Anubis.Client.StateTest do state = new_test_state() from = {self(), make_ref()} - operation = Operation.new(%{method: "test_method"}) + operation = Operation.new(%{method: "test_method", timeout: 30_000}) {request_id, state} = State.add_request_from_operation(state, operation, from) @@ -81,7 +84,7 @@ defmodule Anubis.Client.StateTest do state = new_test_state() from = {self(), make_ref()} - operation = Operation.new(%{method: "test_method"}) + operation = Operation.new(%{method: "test_method", timeout: 30_000}) {request_id, state} = State.add_request_from_operation(state, operation, from) @@ -108,7 +111,7 @@ defmodule Anubis.Client.StateTest do state = new_test_state() from = {self(), make_ref()} - operation = Operation.new(%{method: "test_method"}) + operation = Operation.new(%{method: "test_method", timeout: 30_000}) {request_id, state} = State.add_request_from_operation(state, operation, from) @@ -220,7 +223,7 @@ defmodule Anubis.Client.StateTest do state = new_test_state() from = {self(), make_ref()} - operation = Operation.new(%{method: "test_method"}) + operation = Operation.new(%{method: "test_method", timeout: 30_000}) {request_id, state} = State.add_request_from_operation(state, operation, from) @@ -330,7 +333,8 @@ defmodule Anubis.Client.StateTest do client_info: %{"name" => "TestClient", "version" => "1.0.0"}, capabilities: %{}, protocol_version: "2024-11-05", - transport: %{layer: :fake_transport, name: :fake_name} + transport: %{layer: :fake_transport, name: :fake_name}, + timeout: 30_000 } end end diff --git a/test/anubis/server/transport/sse_test.exs b/test/anubis/server/transport/sse_test.exs index c36995e4..ca7aaafd 100644 --- a/test/anubis/server/transport/sse_test.exs +++ b/test/anubis/server/transport/sse_test.exs @@ -126,7 +126,7 @@ defmodule Anubis.Server.Transport.SSETest do assert_receive :registered, 1000 message = "broadcast message" - assert :ok = SSE.send_message(transport, message) + assert :ok = SSE.send_message(transport, message, timeout: 5000) # Both handlers should receive the message assert_receive {:sse_message, ^message} diff --git a/test/anubis/server/transport/stdio_test.exs b/test/anubis/server/transport/stdio_test.exs index 59695e16..db6979f1 100644 --- a/test/anubis/server/transport/stdio_test.exs +++ b/test/anubis/server/transport/stdio_test.exs @@ -32,7 +32,7 @@ defmodule Anubis.Server.Transport.STDIOTest do message = "test message" assert capture_io(pid, fn -> - assert :ok = STDIO.send_message(pid, message) + assert :ok = STDIO.send_message(pid, message, timeout: 5000) Process.sleep(50) end) =~ "test message" diff --git a/test/anubis/server/transport/streamable_http_test.exs b/test/anubis/server/transport/streamable_http_test.exs index 16601a95..891dbd1e 100644 --- a/test/anubis/server/transport/streamable_http_test.exs +++ b/test/anubis/server/transport/streamable_http_test.exs @@ -123,9 +123,9 @@ defmodule Anubis.Server.Transport.StreamableHTTPTest do end) end - test "send_message/2 works", %{transport: transport} do + test "send_message/3 works", %{transport: transport} do message = "test message" - assert :ok = StreamableHTTP.send_message(transport, message) + assert :ok = StreamableHTTP.send_message(transport, message, timeout: 5000) end test "shutdown/1 gracefully shuts down", %{transport: transport} do diff --git a/test/anubis/transport/sse_test.exs b/test/anubis/transport/sse_test.exs index ab732b94..f2cf6836 100644 --- a/test/anubis/transport/sse_test.exs +++ b/test/anubis/transport/sse_test.exs @@ -131,7 +131,7 @@ defmodule Anubis.Transport.SSETest do {:ok, ping_message} = Message.encode_request(%{"method" => "ping", "params" => %{}}, "1") - assert :ok = SSE.send_message(transport, ping_message) + assert :ok = SSE.send_message(transport, ping_message, timeout: 5000) # Give time for the response to come back Process.sleep(100) @@ -168,7 +168,7 @@ defmodule Anubis.Transport.SSETest do Process.sleep(100) # Try to send a message without having an endpoint - assert {:error, :not_connected} = SSE.send_message(transport, "test message") + assert {:error, :not_connected} = SSE.send_message(transport, "test message", timeout: 5000) # Clean up SSE.shutdown(transport) @@ -222,7 +222,7 @@ defmodule Anubis.Transport.SSETest do # Send a message and check for error response assert {:error, {:http_error, 500, "Internal Server Error"}} = - SSE.send_message(transport, "test message") + SSE.send_message(transport, "test message", timeout: 5000) # Clean up SSE.shutdown(transport) @@ -380,7 +380,7 @@ defmodule Anubis.Transport.SSETest do transport_state = :sys.get_state(transport) assert transport_state.message_url - assert :ok = SSE.send_message(transport, "test message") + assert :ok = SSE.send_message(transport, "test message", timeout: 5000) SSE.shutdown(transport) StubClient.clear_messages() @@ -420,7 +420,7 @@ defmodule Anubis.Transport.SSETest do transport_state = :sys.get_state(transport) assert transport_state.message_url == "#{server_url}/messages/123" - assert :ok = SSE.send_message(transport, "test message") + assert :ok = SSE.send_message(transport, "test message", timeout: 5000) SSE.shutdown(transport) StubClient.clear_messages() @@ -494,7 +494,7 @@ defmodule Anubis.Transport.SSETest do transport_state = :sys.get_state(transport) assert transport_state.message_url == "#{server_url}/messages/123" refute String.contains?(transport_state.message_url, "/mcp/mcp/") - assert :ok = SSE.send_message(transport, "test message") + assert :ok = SSE.send_message(transport, "test message", timeout: 5000) SSE.shutdown(transport) StubClient.clear_messages() diff --git a/test/anubis/transport/stdio_test.exs b/test/anubis/transport/stdio_test.exs index 594a3101..9b388d76 100644 --- a/test/anubis/transport/stdio_test.exs +++ b/test/anubis/transport/stdio_test.exs @@ -57,7 +57,7 @@ defmodule Anubis.Transport.STDIOTest do end test "sends message successfully", %{transport_pid: pid} do - assert :ok = STDIO.send_message(pid, "test message") + assert :ok = STDIO.send_message(pid, "test message", timeout: 5000) end end @@ -80,7 +80,7 @@ defmodule Anubis.Transport.STDIOTest do test "forwards data to client", %{transport_pid: pid} do :ok = StubClient.clear_messages() - STDIO.send_message(pid, "echo test\n") + STDIO.send_message(pid, "echo test\n", timeout: 5000) Process.sleep(100) diff --git a/test/anubis/transport/streamable_http_test.exs b/test/anubis/transport/streamable_http_test.exs index 759b542a..1cfdb9ad 100644 --- a/test/anubis/transport/streamable_http_test.exs +++ b/test/anubis/transport/streamable_http_test.exs @@ -60,7 +60,7 @@ defmodule Anubis.Transport.StreamableHTTPTest do end end - describe "send_message/2" do + describe "send_message/3" do test "sends HTTP POST request with JSON response", %{bypass: bypass} do server_url = "http://localhost:#{bypass.port}" {:ok, stub_client} = StubClient.start_link() @@ -87,7 +87,7 @@ defmodule Anubis.Transport.StreamableHTTPTest do {:ok, ping_message} = Message.encode_request(%{"method" => "ping", "params" => %{}}, "1") - assert :ok = StreamableHTTP.send_message(transport, ping_message) + assert :ok = StreamableHTTP.send_message(transport, ping_message, timeout: 5000) Process.sleep(100) @@ -121,7 +121,7 @@ defmodule Anubis.Transport.StreamableHTTPTest do Process.sleep(100) notification = ~s|{"jsonrpc":"2.0","method":"notifications/initialized"}| - assert :ok = StreamableHTTP.send_message(transport, notification) + assert :ok = StreamableHTTP.send_message(transport, notification, timeout: 5000) Process.sleep(100) @@ -155,7 +155,7 @@ defmodule Anubis.Transport.StreamableHTTPTest do {:ok, ping_message} = Message.encode_request(%{"method" => "ping", "params" => %{}}, "1") - assert :ok = StreamableHTTP.send_message(transport, ping_message) + assert :ok = StreamableHTTP.send_message(transport, ping_message, timeout: 5000) Process.sleep(200) @@ -185,7 +185,7 @@ defmodule Anubis.Transport.StreamableHTTPTest do Process.sleep(100) assert {:error, {:http_error, 500, "Internal Server Error"}} = - StreamableHTTP.send_message(transport, "test message") + StreamableHTTP.send_message(transport, "test message", timeout: 5000) StreamableHTTP.shutdown(transport) StubClient.clear_messages() @@ -211,7 +211,7 @@ defmodule Anubis.Transport.StreamableHTTPTest do Process.sleep(100) assert {:error, {:unsupported_content_type, "text/html"}} = - StreamableHTTP.send_message(transport, "test message") + StreamableHTTP.send_message(transport, "test message", timeout: 5000) StreamableHTTP.shutdown(transport) StubClient.clear_messages() @@ -251,7 +251,7 @@ defmodule Anubis.Transport.StreamableHTTPTest do {:ok, ping_message} = Message.encode_request(%{"method" => "ping", "params" => %{}}, "1") - assert :ok = StreamableHTTP.send_message(transport, ping_message) + assert :ok = StreamableHTTP.send_message(transport, ping_message, timeout: 5000) Process.sleep(100) @@ -305,14 +305,14 @@ defmodule Anubis.Transport.StreamableHTTPTest do {:ok, first_message} = Message.encode_request(%{"method" => "ping", "params" => %{}}, "1") - assert :ok = StreamableHTTP.send_message(transport, first_message) + assert :ok = StreamableHTTP.send_message(transport, first_message, timeout: 5000) Process.sleep(100) {:ok, second_message} = Message.encode_request(%{"method" => "ping", "params" => %{}}, "2") - assert :ok = StreamableHTTP.send_message(transport, second_message) + assert :ok = StreamableHTTP.send_message(transport, second_message, timeout: 5000) Process.sleep(100) @@ -353,7 +353,7 @@ defmodule Anubis.Transport.StreamableHTTPTest do {:ok, ping_message} = Message.encode_request(%{"method" => "ping", "params" => %{}}, "1") - assert :ok = StreamableHTTP.send_message(transport, ping_message) + assert :ok = StreamableHTTP.send_message(transport, ping_message, timeout: 5000) Process.sleep(100) @@ -387,7 +387,7 @@ defmodule Anubis.Transport.StreamableHTTPTest do {:ok, ping_message} = Message.encode_request(%{"method" => "ping", "params" => %{}}, "1") - assert :ok = StreamableHTTP.send_message(transport, ping_message) + assert :ok = StreamableHTTP.send_message(transport, ping_message, timeout: 5000) Process.sleep(100) @@ -414,7 +414,7 @@ defmodule Anubis.Transport.StreamableHTTPTest do Process.sleep(100) assert {:error, _reason} = - StreamableHTTP.send_message(transport, "test message") + StreamableHTTP.send_message(transport, "test message", timeout: 5000) StreamableHTTP.shutdown(transport) StubClient.clear_messages() diff --git a/test/anubis/transport/websocket_test.exs b/test/anubis/transport/websocket_test.exs index 9ccd3e88..b30df314 100644 --- a/test/anubis/transport/websocket_test.exs +++ b/test/anubis/transport/websocket_test.exs @@ -112,7 +112,7 @@ if Code.ensure_loaded?(:gun) do assert Process.alive?(ws_pid) - assert :ok = WebSocket.send_message(ws_pid, test_message) + assert :ok = WebSocket.send_message(ws_pid, test_message, timeout: 5000) assert_receive {:message_sent, ^test_message}, 1000 diff --git a/test/support/mock_transport.ex b/test/support/mock_transport.ex index a7eb018a..d0d1bce3 100644 --- a/test/support/mock_transport.ex +++ b/test/support/mock_transport.ex @@ -6,7 +6,7 @@ defmodule MockTransport do def start_link(_opts), do: {:ok, self()} @impl true - def send_message(_, _), do: :ok + def send_message(_, _, _opts \\ [timeout: 1_000]), do: :ok @impl true def shutdown(_), do: :ok diff --git a/test/support/stub_transport.ex b/test/support/stub_transport.ex index 099ebc33..158c99ef 100644 --- a/test/support/stub_transport.ex +++ b/test/support/stub_transport.ex @@ -80,7 +80,7 @@ defmodule StubTransport do end @impl true - def send_message(transport \\ __MODULE__, message) do + def send_message(transport \\ __MODULE__, message, _opts \\ []) do GenServer.call(transport, {:send_message, message}) end