Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
49 changes: 47 additions & 2 deletions lib/imp/run.ex
Original file line number Diff line number Diff line change
Expand Up @@ -38,7 +38,10 @@ defmodule Imp.Run do
and the snapshot whole however large it is, and with `:max_events` and
`:max_snapshot_bytes` set to `:infinity` nothing is ever evicted. A host that
must keep a complete record of a run sets all three. An oversized event
payload becomes a digest and size marker before sink delivery, and snapshot
payload becomes a digest and size marker before sink delivery. A failed
event keeps a small error marker in its place, with the validated HTTP
status, provider code and retryability when present and when the marker fits,
never the error message or the request and response bodies. Snapshot
eviction adds a `:capture_gap` marker. A sink receives every bounded event;
the snapshot is a bounded recent window.
"""
Expand Down Expand Up @@ -508,15 +511,57 @@ defmodule Imp.Run.Control do
| input: nil,
output: nil,
reasoning: nil,
error: nil,
error: bounded_error(event.error),
metadata:
event.metadata
|> Map.take([:model_call_id])
|> Map.put(:capture, %{truncated: true, original_bytes: bytes, sha256: digest})
}
|> fit_error_summary(max_bytes)
end
end

defp fit_error_summary(%{error: nil} = event, _max_bytes), do: event

defp fit_error_summary(event, max_bytes) do
marker = %{event | error: %{truncated: true}}

cond do
:erlang.external_size(event) <= max_bytes -> event
:erlang.external_size(marker) <= max_bytes -> marker
# The existing capture envelope may itself exceed a very small limit.
# Do not enlarge that envelope when even the failure marker cannot fit.
true -> %{event | error: nil}
end
end

defp bounded_error(nil), do: nil

# Errors can carry an entire provider request. Keep failure distinguishable
# from successful output without retaining messages, bodies, headers or cause.
# Redaction has already run; even these named fields must have bounded shapes.
defp bounded_error(error) when is_map(error) do
error
|> Map.take([:status, :provider_code, :retryable])
|> Enum.reduce(%{truncated: true}, fn
{:status, status}, summary when is_integer(status) and status in 100..599 ->
Map.put(summary, :status, status)

{:retryable, retryable}, summary when is_boolean(retryable) ->
Map.put(summary, :retryable, retryable)

{:provider_code, code}, summary when is_binary(code) and byte_size(code) <= 64 ->
if Regex.match?(~r/\A[A-Za-z0-9_.:-]+\z/, code),
do: Map.put(summary, :provider_code, code),
else: summary

_, summary ->
summary
end)
end

defp bounded_error(_error), do: %{truncated: true}

defp bound_snapshot(state) do
if over?(length(state.events), state.limits.max_events) or
over?(state.snapshot_bytes, state.limits.max_snapshot_bytes) do
Expand Down
103 changes: 103 additions & 0 deletions test/run_observation_test.exs
Original file line number Diff line number Diff line change
Expand Up @@ -137,6 +137,109 @@ defmodule Imp.RunObservationTest do
end)
end

test "oversized provider errors retain status without request or response content" do
{:ok, run} = Imp.Run.start(%Wait{}, %{owner: self()})
assert_receive :waiting

error =
ReqLLM.Error.API.Request.exception(
status: 429,
reason: "private error explanation",
request_body: String.duplicate("x", 215_000) <> "private request content",
response_body: %{"message" => "private response content"},
headers: %{"authorization" => "Bearer private-header-value"}
)
|> Map.put(:provider_code, "rate_limit_exceeded")
|> Map.put(:retryable, true)

Imp.Run.with_context(run.control, fn ->
Imp.Run.emit(:model_response, error: error, metadata: %{model_call_id: "fixture-call"})
end)

event = List.last(Imp.Run.events(run))

assert event.error == %{
truncated: true,
status: 429,
provider_code: "rate_limit_exceeded",
retryable: true
}

assert event.metadata.model_call_id == "fixture-call"
assert event.metadata.capture.truncated
assert event.metadata.capture.original_bytes > 215_000
assert :erlang.external_size(event) < 65_536

serialized = event |> Imp.Run.Event.to_map() |> Jason.encode!()
refute serialized =~ "private"
refute serialized =~ "request_body"
refute serialized =~ "response_body"
refute serialized =~ "authorization"
:ok = Imp.Run.cancel(run)
end

test "error summaries refuse arbitrary fields and unbounded provider codes" do
{:ok, run} = Imp.Run.start(%Wait{}, %{owner: self()}, max_event_bytes: 1000)
assert_receive :waiting

for error <- [
%{status: "429", provider_code: String.duplicate("x", 2000), retryable: "true"},
%{status: 999, provider_code: "sk-test-secret-1234567890", retryable: nil},
%{status: -1, provider_code: "private words", retryable: %{private: "value"}},
%{provider_code: "rate_limit\n"},
{:provider_error, String.duplicate("private", 2000)}
] do
Imp.Run.with_context(run.control, fn ->
Imp.Run.emit(:model_response, error: error, input: String.duplicate("large", 1000))
end)

event = List.last(Imp.Run.events(run))
assert event.error == %{truncated: true}
assert :erlang.external_size(event) <= 1000
end

:ok = Imp.Run.cancel(run)
end

test "an error summary does not enlarge the existing capture envelope past a tight limit" do
for limit <- [300, 400, 500, 600] do
{:ok, run} = Imp.Run.start(%Wait{}, %{owner: self()}, max_event_bytes: limit)
assert_receive :waiting

Imp.Run.with_context(run.control, fn ->
Imp.Run.emit(:model_response,
error: %{status: 429, provider_code: String.duplicate("x", 64), retryable: true},
input: String.duplicate("large", 1000)
)
end)

event = List.last(Imp.Run.events(run))
assert event.input == nil
assert event.metadata.capture.truncated
envelope_bytes = :erlang.external_size(%{event | error: nil})
assert :erlang.external_size(event) <= max(limit, envelope_bytes)
:ok = Imp.Run.cancel(run)
end
end

test "ordinary errors and truncated successful responses keep their existing shape" do
{:ok, run} = Imp.Run.start(%Wait{}, %{owner: self()}, max_event_bytes: 1000)
assert_receive :waiting

Imp.Run.with_context(run.control, fn ->
Imp.Run.emit(:model_response, error: %{status: 429, reason: "short fixture"})
Imp.Run.emit(:model_response, output: String.duplicate("large", 1000))
end)

[_started, ordinary, large] = Imp.Run.events(run)
assert ordinary.error == %{status: 429, reason: "short fixture"}
refute Map.has_key?(ordinary.metadata, :capture)
assert large.error == nil
assert large.output == nil
assert large.metadata.capture.truncated
:ok = Imp.Run.cancel(run)
end

test "task death is recorded once by control even when the task cannot emit" do
{:ok, run} = Imp.Run.start(%Wait{}, %{owner: self()})
assert_receive :waiting
Expand Down
Loading