inference: reverse-proxy admitted Coder turns to the rust coder API

5016a4c0cb1e · AtlantisPleb · · parent eb479349f492

inference: reverse-proxy admitted Coder turns to the rust coder API

When OPENAGENTS_CODER_API_ORIGIN and OPENAGENTS_CODER_API_INTERNAL_TOKEN
are set, Phoenix keeps grant auth and credit, hops the OpenAI body to
rust, streams the SSE back, and meters usage from the rust frames.
Unset, the local adapter path is unchanged.

Deploy story

What this commit did to the running system — joined from the forge receipt chain, the part a commit page elsewhere cannot show.

Not deployed through the forge lane

No push, promotion, build, or deploy receipt references this commit (receipts are scanned over a bounded recent window). Changes shipped by full node replacement carry their proof in the release gate receipt instead.

Changed files

  • modified config/config.exs
  • modified config/runtime.exs
  • added lib/openagents/inference/coder_api_hop.ex
  • modified lib/openagents_web/controllers/inference_proxy_controller.ex
  • added test/openagents/inference/coder_api_hop_test.exs
  • modified test/openagents_web/controllers/inference_proxy_controller_test.exs

Diff

6 files changed, +296 -1

config/config.exs modified +2

@@ -499,6 +499,8 @@ config :openagents,

499 499
  voice_recording_encryption_key: nil,
500 500
  content_encryption_key: nil,
501 501
  inference_proxy_url: nil,
502
  coder_api_origin: nil,
503
  coder_api_internal_token: nil,
502 504
  inference_grant_max_total_tokens: 2_000_000,
503 505
  inference_grant_max_calls: 64,
504 506
  inference_grant_max_cost_microusd: 5_000_000,
config/runtime.exs modified +2

@@ -549,6 +549,8 @@ if config_env() == :prod and runtime_role == :web do

549 549
    vercel_gateway_api_key: optional_text.("AI_GATEWAY_API_KEY"),
550 550
    box_api_key: optional_text.("BOX_API_KEY"),
551 551
    inference_proxy_url: optional_text.("OPENAGENTS_INFERENCE_PROXY_URL"),
552
    coder_api_origin: optional_text.("OPENAGENTS_CODER_API_ORIGIN"),
553
    coder_api_internal_token: optional_text.("OPENAGENTS_CODER_API_INTERNAL_TOKEN"),
552 554
    forge_enabled: forge_enabled,
553 555
    forge_deploy_lane_enabled: forge_deploy_enabled,
554 556
    repository_provisioner_enabled: forge_enabled,
lib/openagents/inference/coder_api_hop.ex added +117

@@ -0,0 +1,117 @@

1
defmodule OpenAgents.Inference.CoderApiHop do
2
  @moduledoc """
3
  Reverse-proxy an already-admitted inference call to the rust coder API.
4
5
  Phoenix keeps grant auth, credit, and metering. The rust process runs the
6
  catalog, simple-flash classifier, and pinned provider stream.
7
  """
8
9
  require Logger
10
11
  @doc "Configured rust origin and shared internal token, or `:local`."
12
  @spec target() :: {:ok, String.t(), String.t()} | :local
13
  def target do
14
    origin = Application.get_env(:openagents, :coder_api_origin)
15
    token = Application.get_env(:openagents, :coder_api_internal_token)
16
17
    if present?(origin) and present?(token) do
18
      {:ok, String.trim_trailing(origin, "/"), token}
19
    else
20
      :local
21
    end
22
  end
23
24
  @doc "POST the OpenAI body to rust and return status, headers, and SSE body."
25
  @spec post(String.t(), String.t(), String.t(), map()) ::
26
          {:ok, pos_integer(), [{String.t(), String.t()}], binary()} | {:error, term()}
27
  def post(origin, internal_token, admitted_model, body) when is_map(body) do
28
    url = origin <> "/api/inference/proxy"
29
30
    case Req.post(url,
31
           json: body,
32
           headers: [
33
             {"authorization", "Bearer " <> internal_token},
34
             {"x-openagents-admitted-model", admitted_model},
35
             {"accept", "text/event-stream"}
36
           ],
37
           receive_timeout: 120_000,
38
           retry: false
39
         ) do
40
      {:ok, %Req.Response{status: status, headers: headers, body: body}} ->
41
        {:ok, status, headers, body_to_binary(body)}
42
43
      {:error, reason} ->
44
        Logger.warning("coder_api_hop_failed reason=#{inspect(reason, limit: 80)}")
45
        {:error, reason}
46
    end
47
  end
48
49
  @doc "OpenAI `usage` object → Phoenix grant usage keys."
50
  @spec usage_from_sse(binary()) :: map()
51
  def usage_from_sse(body) when is_binary(body) do
52
    body
53
    |> String.split("\n\n", trim: true)
54
    |> Enum.reduce(%{}, fn frame, acc ->
55
      payload = String.replace_prefix(frame, "data: ", "")
56
57
      case Jason.decode(payload) do
58
        {:ok, %{"usage" => usage}} when is_map(usage) -> openai_usage(usage)
59
        _ -> acc
60
      end
61
    end)
62
  end
63
64
  def usage_from_sse(_), do: %{}
65
66
  @doc "Header value rust attributes as the served model."
67
  @spec served_model(term()) :: String.t() | nil
68
  def served_model(headers) when is_map(headers) do
69
    case Map.get(headers, "x-openagents-model") do
70
      [value | _] when is_binary(value) -> value
71
      value when is_binary(value) -> value
72
      _ -> nil
73
    end
74
  end
75
76
  def served_model(headers) when is_list(headers) do
77
    Enum.find_value(headers, fn
78
      {name, value} ->
79
        if String.downcase(to_string(name)) == "x-openagents-model" do
80
          header_value(value)
81
        end
82
83
      _ ->
84
        nil
85
    end)
86
  end
87
88
  def served_model(_), do: nil
89
90
  defp openai_usage(usage) do
91
    input = integer(usage["prompt_tokens"] || usage["input_tokens"])
92
    output = integer(usage["completion_tokens"] || usage["output_tokens"])
93
    total = integer(usage["total_tokens"])
94
95
    %{
96
      "input_tokens" => input,
97
      "output_tokens" => output,
98
      "total_tokens" => if(total > 0, do: total, else: input + output)
99
    }
100
    |> Enum.reject(fn {_k, v} -> v == 0 end)
101
    |> Map.new()
102
  end
103
104
  defp body_to_binary(body) when is_binary(body), do: body
105
  defp body_to_binary(%Req.Response.Async{} = async), do: Enum.into(async, "")
106
  defp body_to_binary(_), do: ""
107
108
  defp header_value([value | _]) when is_binary(value), do: value
109
  defp header_value(value) when is_binary(value), do: value
110
  defp header_value(_), do: nil
111
112
  defp present?(value) when is_binary(value), do: String.trim(value) != ""
113
  defp present?(_), do: false
114
115
  defp integer(value) when is_integer(value) and value >= 0, do: value
116
  defp integer(_), do: 0
117
end
lib/openagents_web/controllers/inference_proxy_controller.ex modified +66 -1

@@ -64,12 +64,77 @@ defmodule OpenAgentsWeb.InferenceProxyController do

64 64
         :ok <- serving(model),
65 65
         :ok <- requested_model(model, conn.body_params),
66 66
         {:ok, request} <- build_request(model, conn.body_params) do
67
      run(conn, grant, model, request)
67
      case OpenAgents.Inference.CoderApiHop.target() do
68
        {:ok, origin, token} -> hop(conn, grant, model, request, origin, token)
69
        :local -> run(conn, grant, model, request)
70
      end
68 71
    else
69 72
      {:error, reason} -> refuse(conn, reason)
70 73
    end
71 74
  end
72 75
76
  defp hop(conn, grant, model, request, origin, token) do
77
    selection = selection_properties(grant, model, request, conn.body_params)
78
    Analytics.capture("inference_model_selected", analytics_distinct_id(grant), selection)
79
80
    case OpenAgents.Inference.CoderApiHop.post(origin, token, model.id, conn.body_params) do
81
      {:ok, 200, headers, body} ->
82
        served_name = OpenAgents.Inference.CoderApiHop.served_model(headers) || model.id
83
        served = if served_name == model.id, do: :requested, else: served_name
84
        usage = OpenAgents.Inference.CoderApiHop.usage_from_sse(body)
85
        _ = meter(grant, usage, served)
86
        record_health(model, served)
87
        label = model_label(model, served)
88
89
        Analytics.capture(
90
          "inference_model_served",
91
          analytics_distinct_id(grant),
92
          Map.merge(selection, %{
93
            "served_model" => label,
94
            "served_model_disclosed" => served != :unresolved,
95
            "outcome" => "served",
96
            "usage_reported" => usage != %{},
97
            "coder_api_hop" => true
98
          })
99
        )
100
101
        conn
102
        |> put_resp_content_type("text/event-stream")
103
        |> put_resp_header("cache-control", "no-store")
104
        |> put_resp_header("x-openagents-model", label)
105
        |> send_resp(200, body)
106
107
      {:ok, status, _headers, _body} ->
108
        Logger.warning("coder_api_hop_refused status=#{status}")
109
110
        Analytics.capture(
111
          "inference_model_failed",
112
          analytics_distinct_id(grant),
113
          Map.merge(selection, %{
114
            "outcome" => "provider_failed",
115
            "reason_code" => "coder_api_hop",
116
            "upstream_status" => status,
117
            "usage_reported" => false
118
          })
119
        )
120
121
        refuse(conn, {:provider_failed, "coder_api_hop", status})
122
123
      {:error, _reason} ->
124
        Analytics.capture(
125
          "inference_model_failed",
126
          analytics_distinct_id(grant),
127
          Map.merge(selection, %{
128
            "outcome" => "provider_failed",
129
            "reason_code" => "coder_api_hop",
130
            "usage_reported" => false
131
          })
132
        )
133
134
        refuse(conn, {:provider_failed, "coder_api_hop", nil})
135
    end
136
  end
137
73 138
  # ── request assembly ────────────────────────────────────────────────────
74 139
75 140
  # The grant's model names the provider and the string that provider is called
test/openagents/inference/coder_api_hop_test.exs added +40

@@ -0,0 +1,40 @@

1
defmodule OpenAgents.Inference.CoderApiHopTest do
2
  use ExUnit.Case, async: true
3
4
  alias OpenAgents.Inference.CoderApiHop
5
6
  test "usage_from_sse reads OpenAI prompt and completion tokens" do
7
    body = """
8
    data: {"choices":[{"delta":{"content":"hi"}}],"model":"gemini-3.7-flash"}
9
10
    data: {"choices":[],"usage":{"prompt_tokens":10,"completion_tokens":2,"total_tokens":12}}
11
12
    data: [DONE]
13
14
    """
15
16
    assert CoderApiHop.usage_from_sse(body) == %{
17
             "input_tokens" => 10,
18
             "output_tokens" => 2,
19
             "total_tokens" => 12
20
           }
21
  end
22
23
  test "served_model reads the rust attribution header" do
24
    assert CoderApiHop.served_model(%{"x-openagents-model" => ["gemini-3.7-flash"]}) ==
25
             "gemini-3.7-flash"
26
27
    assert CoderApiHop.served_model([{"x-openagents-model", "glm-5.3-flash"}]) == "glm-5.3-flash"
28
    assert CoderApiHop.served_model(%{}) == nil
29
  end
30
31
  test "target is local when origin or token is missing" do
32
    old_origin = Application.get_env(:openagents, :coder_api_origin)
33
    old_token = Application.get_env(:openagents, :coder_api_internal_token)
34
    Application.put_env(:openagents, :coder_api_origin, nil)
35
    Application.put_env(:openagents, :coder_api_internal_token, nil)
36
    assert CoderApiHop.target() == :local
37
    Application.put_env(:openagents, :coder_api_origin, old_origin)
38
    Application.put_env(:openagents, :coder_api_internal_token, old_token)
39
  end
40
end
test/openagents_web/controllers/inference_proxy_controller_test.exs modified +69

@@ -716,4 +716,73 @@ defmodule OpenAgentsWeb.InferenceProxyControllerTest do

716 716
      end
717 717
    end
718 718
  end
719
720
  test "a coder-api hop streams rust SSE, attributes the served model, and meters", %{conn: conn} do
721
    parent = self()
722
723
    {:ok, pid} =
724
      Bandit.start_link(
725
        plug: {__MODULE__.CoderHopStub, parent},
726
        scheme: :http,
727
        port: 0,
728
        thousand_island_options: [num_acceptors: 1]
729
      )
730
731
    {:ok, {_address, port}} = ThousandIsland.listener_info(pid)
732
    origin = "http://127.0.0.1:#{port}"
733
734
    old_origin = Application.get_env(:openagents, :coder_api_origin)
735
    old_token = Application.get_env(:openagents, :coder_api_internal_token)
736
    Application.put_env(:openagents, :coder_api_origin, origin)
737
    Application.put_env(:openagents, :coder_api_internal_token, "hop-secret")
738
739
    on_exit(fn ->
740
      Application.put_env(:openagents, :coder_api_origin, old_origin)
741
      Application.put_env(:openagents, :coder_api_internal_token, old_token)
742
      Process.exit(pid, :normal)
743
    end)
744
745
    %{grant: grant, token: token} = grant("hop")
746
747
    conn =
748
      post_chat(conn, token, %{
749
        "model" => Models.default_id(),
750
        "messages" => [%{"role" => "user", "content" => "hey"}]
751
      })
752
753
    assert conn.status == 200
754
    assert get_resp_header(conn, "x-openagents-model") == ["gemini-3.7-flash"]
755
    assert conn.resp_body =~ "Ready"
756
    assert_receive {:coder_hop, ["Bearer hop-secret"], [admitted]}
757
    assert admitted == Models.default_id()
758
759
    metered = Repo.get(Grant, grant.id)
760
    assert metered.call_count == 1
761
    assert metered.usage["input_tokens"] == 10
762
    assert metered.usage["output_tokens"] == 2
763
  end
764
end
765
766
defmodule OpenAgentsWeb.InferenceProxyControllerTest.CoderHopStub do
767
  @behaviour Plug
768
769
  def init(parent), do: parent
770
771
  def call(conn, parent) do
772
    send(
773
      parent,
774
      {:coder_hop, Plug.Conn.get_req_header(conn, "authorization"),
775
       Plug.Conn.get_req_header(conn, "x-openagents-admitted-model")}
776
    )
777
778
    body =
779
      "data: {\"choices\":[{\"delta\":{\"content\":\"Ready\"}}],\"model\":\"gemini-3.7-flash\"}\n\n" <>
780
        "data: {\"choices\":[],\"usage\":{\"prompt_tokens\":10,\"completion_tokens\":2,\"total_tokens\":12}}\n\n" <>
781
        "data: [DONE]\n\n"
782
783
    conn
784
    |> Plug.Conn.put_resp_header("x-openagents-model", "gemini-3.7-flash")
785
    |> Plug.Conn.put_resp_content_type("text/event-stream")
786
    |> Plug.Conn.send_resp(200, body)
787
  end
719 788
end

This page updates live while a promote is in flight · changelog