Skip to content
2 changes: 1 addition & 1 deletion .dialyzer_ignore.exs
Original file line number Diff line number Diff line change
Expand Up @@ -55,7 +55,7 @@
# defensive guard: ReqLLM.Response types provider_meta as map() with a %{}
# default, but the struct does not enforce it (a caller can build one with
# nil), and ReqLLM's own OpenTelemetry attributes guard it with is_map/1.
{"lib/imp/clients/req_llm.ex", :guard_fail, 1902},
{"lib/imp/clients/req_llm.ex", :guard_fail, 1917},
# defensive error clause on an always-ok internal call
{"lib/imp/clients/training.ex", :pattern_match, {1215, 13}},
# defensive error clause on an always-ok internal call
Expand Down
30 changes: 29 additions & 1 deletion CHANGELOG.md
Original file line number Diff line number Diff line change
Expand Up @@ -4,7 +4,7 @@ User-visible changes to Imp are recorded here.

## Unreleased

### Fixed
### Changed

- Breaking: `Imp.Core.LMResponse.cost`, and the `:cost` on a
`:model_response` event, is what the provider reported charging, as its
Expand All @@ -23,6 +23,34 @@ User-visible changes to Imp are recorded here.
unknown charge, not a free one; a host that wants the old number for calls
with no reported charge reads `estimated_cost` for them, knowing it is an
estimate.

### Fixed

- A call streamed through `Imp.Clients.ReqLLM` records the usage
and the cost the provider reported at the end of the stream. ReqLLM's
stream is a `Stream.resource`, which reports its end as halted rather than
done, and Imp ended such a stream without its terminal event, so every
streamed call was recorded with no usage and no cost. The stream now ends
with exactly one terminal event, `done: true` with the provider's usage
(including its `"cost"`), model and finish reason, taking the finish
reason, and usage no chunk reported, from ReqLLM's metadata handle.
- A stream that did not complete ends in `{:error, %Imp.LMError{}}`, not in
a completion of whatever text arrived first: one that carries a provider
error, finishes with reason `:error` or `:cancelled`, or is incomplete,
its body ending with no finish and no `[DONE]` (reason
`{:stream_finished, :incomplete}`). The error event carries the usage and
other metadata that arrived before it.
- A streamed call that fails after the provider reported usage records that
usage and cost on its failed `:model_response` event and in
`Imp.Usage`, since the provider may have charged for it. The caller still
receives `{:error, reason}`.
- A streamed call to a client built from an inline spec map, with atom or
string keys, or a `{provider, opts}` tuple no longer fails with
`{:lm_stream_failed, "protocol String.Chars not implemented ..."}`. A
streamed call records the model the provider reported, as a non-streamed
call does, and otherwise the configured model id (`gpt-test`, where it
recorded the whole `openai:gpt-test`), so its `Imp.Usage` key changes the
same way.
- `:reasoning_effort` accepts `max`, when an LM is built, on a call and in a
saved program. Imp's accepted efforts are read from ReqLLM's own
`reasoning_effort` option, so they are every effort ReqLLM accepts, on every
Expand Down
170 changes: 137 additions & 33 deletions lib/imp/clients/req_llm.ex
Original file line number Diff line number Diff line change
Expand Up @@ -25,6 +25,15 @@ defmodule Imp.Clients.ReqLLM do
relays an upstream refusal, is returned as that error, never as an empty
completion.

`stream/3` ends with exactly one terminal event. A provider stream that runs
to its end closes with `done: true` and the metadata the provider reported
along the way: usage (with any `"cost"` the provider charged), model and
finish reason. One that raises, carries a provider error, or finishes with
reason `:error` or `:cancelled` closes with `{:error, %Imp.LMError{}}`
instead, never a completion, and that event carries the metadata that
arrived before it stopped. A consumer that stops early receives no terminal
event, and the provider stream is cancelled.

`:reasoning_effort` is the one reasoning option, on the client or on a call.
It takes any value of ReqLLM's own `reasoning_effort` option, such as `high`,
`xhigh`, `max` or `default`, as an atom or a string, on every provider;
Expand Down Expand Up @@ -392,26 +401,28 @@ defmodule Imp.Clients.ReqLLM do
# says nothing has declined to answer (`Imp.Predict.ReActV2`), so a refused
# request would be recorded as a choice. It is the failed request it reports,
# in the shape ReqLLM gives an HTTP error.
defp relayed_error(%ReqLLM.Response{provider_meta: %{} = meta}) do
case Map.get(meta, "error") || Map.get(meta, :error) do
%{} = error ->
code = Map.get(error, "code") || Map.get(error, :code)

%ReqLLM.Error.API.Request{
reason: Map.get(error, "message") || Map.get(error, :message) || inspect(error),
status: if(is_integer(code), do: code),
response_body: %{"error" => error}
}
defp relayed_error(%ReqLLM.Response{provider_meta: %{} = meta}),
do: provider_error(Map.get(meta, "error") || Map.get(meta, :error))

message when is_binary(message) and message != "" ->
%ReqLLM.Error.API.Request{reason: message, response_body: %{"error" => message}}
defp relayed_error(_response), do: nil

_none ->
nil
end
# A provider's error object or message, as ReqLLM carries it in a response's
# `provider_meta` or a stream's metadata, in the shape ReqLLM gives an HTTP
# error; `nil` when there is none.
defp provider_error(%{} = error) do
code = Map.get(error, "code") || Map.get(error, :code)

%ReqLLM.Error.API.Request{
reason: Map.get(error, "message") || Map.get(error, :message) || inspect(error),
status: if(is_integer(code), do: code),
response_body: %{"error" => error}
}
end

defp relayed_error(_response), do: nil
defp provider_error(message) when is_binary(message) and message != "",
do: %ReqLLM.Error.API.Request{reason: message, response_body: %{"error" => message}}

defp provider_error(_none), do: nil

# Every failed request becomes one `Imp.LMError`, classified here, where the
# provider library's error shapes are known, so no caller has to know them.
Expand Down Expand Up @@ -1832,7 +1843,11 @@ defmodule Imp.Clients.ReqLLM do
defp model_id(%{"id" => id}) when is_binary(id), do: id
defp model_id(%{model: id}) when is_binary(id), do: id
defp model_id(%{"model" => id}) when is_binary(id), do: id
defp model_id({_provider, id}) when is_binary(id), do: id
# ReqLLM's two-element tuple is `{provider, opts}`, naming the model in
# `:id` or `:model` (`ReqLLM.model/1`).
defp model_id({_provider, opts}) when is_list(opts),
do: to_string(opts[:id] || opts[:model] || "")

defp model_id({_provider, id, _opts}) when is_binary(id), do: id

defp model_id(model) when is_binary(model) do
Expand Down Expand Up @@ -1914,16 +1929,32 @@ defmodule Imp.Clients.ReqLLM do
}
end

defp provider_name(%{provider: provider}), do: to_string(provider)
defp provider_name(model_spec), do: model_spec |> model_identity() |> elem(0)

defp provider_name(model_spec) when is_binary(model_spec) do
case String.split(model_spec, ":", parts: 2) do
@doc false
# The provider, as a string, and the model id of any model shape
# `Imp.req_llm/2` accepts: a `"provider:model"` string, a
# `{provider, opts}` or `{provider, model, opts}` tuple, or a spec map
# with atom or string keys. The provider is `nil` when the shape names none.
@spec model_identity(term()) :: {String.t() | nil, String.t()}
def model_identity(model), do: {provider_label(model), model_id(model)}

defp provider_label(%{provider: provider}) when not is_nil(provider), do: to_string(provider)

defp provider_label(%{"provider" => provider}) when not is_nil(provider),
do: to_string(provider)

defp provider_label({provider, _model}) when is_atom(provider), do: to_string(provider)
defp provider_label({provider, _model, _opts}) when is_atom(provider), do: to_string(provider)

defp provider_label(model) when is_binary(model) do
case String.split(model, ":", parts: 2) do
[provider, _model] -> provider
_other -> nil
end
end

defp provider_name(_model_spec), do: nil
defp provider_label(_model), do: nil

defp sanitize_usage(nil), do: nil
defp sanitize_usage(usage) when is_map(usage), do: sanitize_usage_value(usage)
Expand Down Expand Up @@ -2134,23 +2165,89 @@ defmodule Imp.Clients.ReqLLM do
metadata: metadata
}}

{:done, _acc} ->
{[
%Imp.Streaming.Messages.StreamResponse{
done: true,
metadata: state.metadata
}
], %{state | continuation: nil, started?: true, completed?: true}}

{:halted, _acc} ->
{:halt, %{state | continuation: nil, started?: true, completed?: true}}
# The reducer only ever suspends and is only ever resumed with
# `{:cont, _}`, so neither result means a consumer stopped early: a list
# that runs out reports `{:done, _}`, and a `Stream.resource` (ReqLLM's
# stream is one) that runs out after a suspension reports `{:halted, _}`.
# Both are the provider stream's end.
{finished, _acc} when finished in [:done, :halted] ->
finish_stream(%{state | continuation: nil, started?: true})
end
rescue
error -> stream_failure(state, error)
catch
kind, reason -> stream_failure(state, {kind, reason})
end

# The end of the provider stream closes with one terminal event carrying
# the metadata accumulated along the way. A stream whose metadata reports an
# error, or a finish reason of `:error`, `:cancelled` or `:incomplete`, did
# not complete, and ends as a failure. ReqLLM's own event projection
# (`ReqLLM.StreamResponse.events/1`) ends the first three the same way; it
# reports `:incomplete` as the reason of a `:finish` event, which Imp does
# not count as a completion.
#
# The chunks carry what the provider sent; ReqLLM's metadata handle carries
# what ReqLLM concluded about the stream as a whole. Its finish reason is
# `:incomplete` when the body ended with no termination event, a stream
# cut short that the chunks alone cannot tell from a finished one, and its
# usage stands in when no chunk reported any.
defp finish_stream(state) do
state = %{state | metadata: merge_handle_metadata(state.metadata, state.response)}

case stream_end_error(state.metadata) do
nil ->
{[%Imp.Streaming.Messages.StreamResponse{done: true, metadata: state.metadata}],
%{state | completed?: true}}

error ->
stream_failure(state, error)
end
end

# The provider stream has ended, so ReqLLM's collection is finishing too;
# the wait is bounded so a handle that never answers cannot hold the
# terminal event, and a handle that fails or has stopped adds nothing.
@metadata_handle_timeout 5_000

defp merge_handle_metadata(metadata, %ReqLLM.StreamResponse{metadata_handle: handle})
when is_pid(handle) do
handle_metadata =
try do
ReqLLM.StreamResponse.MetadataHandle.await(handle, @metadata_handle_timeout)
rescue
_error -> %{}
catch
:exit, _reason -> %{}
end

metadata
|> put_handle_value(:finish_reason, handle_metadata, :always)
|> put_handle_value(:usage, handle_metadata, :when_missing)
end

defp merge_handle_metadata(metadata, _response), do: metadata

defp put_handle_value(metadata, key, handle_metadata, rule) do
case {Map.get(handle_metadata, key), rule, Map.get(metadata, key)} do
{nil, _rule, _current} -> metadata
{value, :always, _current} -> Map.put(metadata, key, value)
{value, :when_missing, nil} -> Map.put(metadata, key, value)
{_value, :when_missing, _current} -> metadata
end
end

defp stream_end_error(metadata) do
with nil <- provider_error(Map.get(metadata, :error) || Map.get(metadata, "error")) do
case Map.get(metadata, :finish_reason) || Map.get(metadata, "finish_reason") do
reason when reason in [:error, "error"] -> {:stream_finished, :error}
reason when reason in [:cancelled, "cancelled"] -> {:stream_finished, :cancelled}
reason when reason in [:incomplete, "incomplete"] -> {:stream_finished, :incomplete}
_other -> nil
end
end
end

defp suspend_stream(stream) do
Enumerable.reduce(stream, {:cont, nil}, fn chunk, _acc -> {:suspend, chunk} end)
end
Expand All @@ -2166,11 +2263,18 @@ defmodule Imp.Clients.ReqLLM do

# The stream had opened, so the request reached the provider: sending it
# again may be billed again, and repeats chunks the caller already has.
# The terminal event is an error, never a completion, and carries whatever
# metadata (usage, cost, finish reason) arrived before the stream broke.
defp stream_failure(state, error) do
reason = lm_error(error, true, state.provider)

{[%Imp.Streaming.Messages.StreamResponse{chunk: {:error, reason}, done: true}],
%{state | completed?: true, failed?: true}}
{[
%Imp.Streaming.Messages.StreamResponse{
chunk: {:error, reason},
done: true,
metadata: state.metadata
}
], %{state | completed?: true, failed?: true}}
end

defp cleanup_stream(state, lm) do
Expand Down
5 changes: 5 additions & 0 deletions lib/imp/exceptions.ex
Original file line number Diff line number Diff line change
Expand Up @@ -42,6 +42,11 @@ defmodule Imp.LMError do
request because its input is longer than the model accepts. Sending it
again unchanged will fail again; a shorter input may not.

An error the provider sends inside a stream arrives without its code,
because ReqLLM's stream decoder keeps only its message, so it has `:status`
`nil` and is retryable as a failed stream, where the same error in a
non-streamed response may carry a status that says otherwise.

`:reason` keeps the provider library's error unchanged for diagnostics.
`Imp.Errors.retryable?/1` and `Imp.Errors.context_window_exceeded?/1` read
these fields through the wrappers Imp puts around an LM error.
Expand Down
Loading
Loading