phoenix_tracker_available? = Code.ensure_loaded?(Phoenix.Tracker) and Code.ensure_loaded?(Phoenix.PubSub) defmodule Phantom.Tracker do @moduledoc """ Track open streams so that notifications and requests can be sent to clients. Add to your supervision tree: {Phoenix.PubSub, name: MyApp.PubSub}, {Phantom.Tracker, [name: Phantom.Tracker, pubsub_server: MyApp.PubSub]}, For example, a request may need to elicit more input from the client, so the first request stream will remain open, and the notification stream will send and new request to the client, and the client will POST its response. The new response connection will notify the first request connection with the result and the tool can continue with the elicited information. See `m:Phantom#module-persistent-streams` section for more information. """ if phoenix_tracker_available? do Code.eval_quoted(quote(do: use(Phoenix.Tracker)), [], __ENV__) end @available phoenix_tracker_available? @sessions "phantom:sessions" @requests "phantom:requests" @resources "phantom:resources" def available?, do: @available def resource_subscription_topic, do: @resources def requests_topic, do: @requests def sessions_topic, do: @sessions @doc false if @available do def start_link(opts) do opts = Keyword.merge([name: __MODULE__], opts) Phoenix.Tracker.start_link(__MODULE__, opts, opts) end else def start_link(_opts), do: :error end @doc false if @available do def init(opts) do server = Keyword.fetch!(opts, :pubsub_server) {:ok, %{pubsub_server: server, node_name: Phoenix.PubSub.node_name(server)}} end else def init(_opts), do: :ignore end @doc "Track a request PID" if @available do def track_request(pid, request_id, meta \\ %{}) do Phoenix.Tracker.track(__MODULE__, pid, @requests, request_id, Map.put(meta, :pid, pid)) rescue _ -> {:error, :tracker_not_in_supervision_tree} end else def track_request(_pid, _request_id, _meta \\ %{}), do: {:error, :not_available} end @doc """ Claim an in-flight JSON-RPC request for a session to prevent double-dispatch when the same `tools/call` lands on two nodes (e.g. proxy retry). Returns `:ok` if the claim is held by this node alone, `:duplicate` if another node has already claimed it. Stored on the `#{@requests}` topic alongside server-initiated request tracking. The `{:in_flight, session_id, request_id}` tuple key never collides with elicitation keys (binaries), so both kinds of tracking share one replication stream. The check is best-effort against Phoenix.Tracker replication lag: two nodes racing can both see `:ok`. Callers should treat this as a hint and combine with deterministic request IDs for defence in depth. """ if @available do def track_in_flight(session_id, request_id) do key = in_flight_key(session_id, request_id) case Phoenix.Tracker.get_by_key(__MODULE__, @requests, key) do [] -> case Phoenix.Tracker.track(__MODULE__, self(), @requests, key, %{ type: :in_flight, pid: self() }) do {:ok, _} -> :ok _ -> :duplicate end _entries -> :duplicate end rescue _ -> :ok end else def track_in_flight(_session_id, _request_id), do: :ok end @doc "Release an in-flight claim taken by `track_in_flight/2`." if @available do def untrack_in_flight(session_id, request_id) do Phoenix.Tracker.untrack( __MODULE__, self(), @requests, in_flight_key(session_id, request_id) ) :ok rescue _ -> :ok end else def untrack_in_flight(_session_id, _request_id), do: :ok end # Tuple keys never collide with the binary/number keys that # `track_request/3` uses for server-initiated requests, so # both kinds of tracking can share the `#{@requests}` topic # and its single replication stream. defp in_flight_key(session_id, request_id), do: {:in_flight, session_id, request_id} @doc "Track a session PID" if @available do require Logger def track_session(pid, session_id, meta \\ %{}) do Phoenix.Tracker.track(__MODULE__, pid, @sessions, session_id, Map.put(meta, :pid, pid)) rescue _ -> if not warned?() do Logger.warning( "Phoenix.PubSub is available, but Phantom.Tracker is not in the supervision tree. Please add this to your supervision tree: `{Phantom.Tracker, [name: Phantom.Tracker, pubsub_server: MyApp.PubSub]}`" ) end {:error, :tracker_not_in_supervision_tree} end defp warned? do :persistent_term.get({Phantom.Tracker, :warned}, false) == true || :persistent_term.put({Phantom.Tracker, :warned}, true) != :ok end else def track_session(_pid, _session_id, _meta \\ %{}), do: {:error, :not_available} end @doc "Update metadata for a tracked session. Falls back to process dictionary when Phoenix.Tracker is not tracking (e.g. stdio transport). " if @available do def update_session_meta(session_id, meta) do case Phoenix.Tracker.get_by_key(__MODULE__, @sessions, session_id) do [{_key, %{pid: pid} = existing_meta} | _] -> Phoenix.Tracker.update( __MODULE__, pid, @sessions, session_id, Map.merge(existing_meta, meta) ) _ -> update_local_session_meta(meta) end rescue _ -> update_local_session_meta(meta) end else def update_session_meta(_session_id, meta), do: update_local_session_meta(meta) end @doc "Get metadata for a tracked session. Falls back to process dictionary when Phoenix.Tracker is not tracking. " if @available do def get_session_meta(session_id) do case Phoenix.Tracker.get_by_key(__MODULE__, @sessions, session_id) do [] -> Process.get(:phantom_session_meta, %{}) entries -> # Merge all entries — the initialize entry has client_capabilities, # a GET entry may not. Merging ensures we find capabilities regardless # of which entry appears first. Enum.reduce(entries, %{}, fn {_key, meta}, acc -> Map.merge(acc, meta) end) end rescue _ -> Process.get(:phantom_session_meta, %{}) end else def get_session_meta(_session_id), do: Process.get(:phantom_session_meta, %{}) end defp update_local_session_meta(meta) do existing = Process.get(:phantom_session_meta, %{}) Process.put(:phantom_session_meta, Map.merge(existing, meta)) :ok end @doc "Return all tracked streams for a given session ID (across all nodes)" if @available do def list_session_streams(session_id) do Phoenix.Tracker.get_by_key(__MODULE__, @sessions, session_id) rescue _ -> [] end else def list_session_streams(_session_id), do: [] end @doc "Return a list of all open sessions" if @available do def list_sessions do Phoenix.Tracker.list(__MODULE__, @sessions) rescue _ -> [] end else def list_sessions, do: [] end @doc "Return a list of all open requests" if @available do def list_requests do Phoenix.Tracker.list(__MODULE__, @requests) rescue _ -> [] end else def list_requests, do: [] end @doc "Return a list of all listening for resources" if @available do def list_resource_listeners do Phoenix.Tracker.list(__MODULE__, @resources) rescue _ -> [] end else def list_resource_listeners, do: [] end @doc "Fetch the PID of the open request by ID" def get_request(%Phantom.Request{id: request_id}), do: get_request(request_id) if @available do def get_request(request_id) do case Phoenix.Tracker.get_by_key(__MODULE__, @requests, request_id) do [{_key, %{pid: pid}} | _] -> pid _ -> nil end rescue _ -> nil end else def get_request(_request_id), do: nil end @doc "Fetch the full metadata map of the open request by ID" def get_request_meta(%Phantom.Request{id: request_id}), do: get_request_meta(request_id) if @available do def get_request_meta(request_id) do case Phoenix.Tracker.get_by_key(__MODULE__, @requests, request_id) do [{_key, meta} | _] -> meta _ -> nil end rescue _ -> nil end else def get_request_meta(_request_id), do: nil end @doc "Fetch the PID of the open session by ID" def get_session(%Phantom.Session{id: session_id}), do: get_session(session_id) if @available do def get_session(session_id) do case Phoenix.Tracker.get_by_key(__MODULE__, @sessions, session_id) do [{_key, %{pid: pid}} | _] -> pid _ -> nil end rescue _ -> nil end else def get_session(_session_id), do: nil end @doc "Untrack the processe for everything" if @available do def untrack(pid), do: Phoenix.Tracker.untrack(__MODULE__, pid) else def untrack(_pid), do: :ok end @doc "Untrack any processes for the session" def untrack_session(%Phantom.Session{id: session_id}), do: untrack_session(session_id) if @available do def untrack_session(session_id) do tracked = try do Phoenix.Tracker.get_by_key(__MODULE__, @sessions, session_id) rescue _ -> [] end Enum.each(tracked, fn {_key, %{pid: pid}} -> Phoenix.Tracker.untrack(__MODULE__, pid) end) :ok rescue _ -> :ok end else def untrack_session(_session_id), do: :ok end @doc "Untrack any processes for the request" def untrack_request(%Phantom.Request{id: request_id}), do: untrack_request(request_id) if @available do def untrack_request(request_id) do tracked = try do Phoenix.Tracker.get_by_key(__MODULE__, @requests, request_id) rescue _ -> [] end Enum.each(tracked, fn {_key, %{pid: pid}} -> Phoenix.Tracker.untrack(__MODULE__, pid) end) :ok rescue _ -> :ok end else def untrack_request(_request_id), do: :ok end @doc "Subscribe the process to resource notifications from the PubSub on topic #{inspect(@resources)}" if @available do def subscribe_resource(uri) do Phoenix.Tracker.track(__MODULE__, self(), @resources, uri, %{pid: self()}) rescue _ -> {:error, :tracker_not_in_supervision_tree} end else def subscribe_resource(_uri), do: {:error, :not_available} end @doc "Unsubscribe the process to resource notifications from the PubSub on topic #{inspect(@resources)}" if @available do def unsubscribe_resource(uri) do Phoenix.Tracker.untrack(__MODULE__, self(), @resources, uri) rescue _ -> :ok end else def unsubscribe_resource(_uri), do: {:error, :not_available} end @doc "Notify any listening MCP sessions that the resource has updated" if @available do def notify_resource_updated(uri) do tracked = try do Phoenix.Tracker.get_by_key(__MODULE__, @resources, uri) rescue _ -> [] end {:ok, Enum.count(tracked, fn {_key, %{pid: pid}} -> GenServer.cast(pid, {:resource_updated, uri}) end)} rescue _ -> {:error, :tracker_not_in_supervision_tree} end else def notify_resource_updated(_), do: {:ok, 0} end @doc "Notify any listening MCP sessions that the list of tools has updated" if @available do def notify_tool_list do {:ok, Enum.count(list_sessions(), fn {session_id, _} -> if pid = get_session(session_id), do: GenServer.cast(pid, :tools_updated) end)} end else def notify_tool_list, do: {:ok, 0} end @doc "Notify any listening MCP sessions that the list of prompts has updated" if @available do def notify_prompt_list do {:ok, Enum.count(list_sessions(), fn {session_id, _} -> if pid = get_session(session_id), do: GenServer.cast(pid, :prompts_updated) end)} end else def notify_prompt_list, do: {:ok, 0} end @doc "Notify any listening MCP sessions that the list of prompts has updated" if @available do def notify_resource_list do {:ok, Enum.count(list_sessions(), fn {session_id, _} -> if pid = get_session(session_id), do: GenServer.cast(pid, :resources_updated) end)} end else def notify_resource_list, do: {:ok, 0} end @doc "Notify tracked sessions that a URL elicitation has completed" if @available do def notify_elicitation_complete(elicitation_id) do tracked = try do Phoenix.Tracker.get_by_key(__MODULE__, @requests, elicitation_id) rescue _ -> [] end Enum.each(tracked, fn {_key, %{type: :elicitation, pid: pid}} -> GenServer.cast(pid, {:send, Phantom.Request.elicitation_complete(elicitation_id)}) Phoenix.Tracker.untrack(__MODULE__, pid, @requests, elicitation_id) end) {:ok, length(tracked)} end else def notify_elicitation_complete(_elicitation_id), do: {:ok, 0} end @doc false if @available do def handle_diff(diff, state) do for {topic, {joins, leaves}} <- diff do for {key, meta} <- joins do msg = {:join, key, meta} Phoenix.PubSub.direct_broadcast!(state.node_name, state.pubsub_server, topic, msg) end for {key, meta} <- leaves do msg = {:leave, key, meta} Phoenix.PubSub.direct_broadcast!(state.node_name, state.pubsub_server, topic, msg) end end {:ok, state} end else def handle_diff(_diff, state), do: {:ok, state} end end