Skip to content
Open
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 @@ -25,7 +25,10 @@ defmodule Imp.Run do
Capture defaults to 64 KiB per event and a 512-event / 4 MiB snapshot. Configure
`:max_event_bytes`, `:max_events`, and `:max_snapshot_bytes` at start. Oversized
event payloads become explicit digest/size markers before sink delivery;
snapshot eviction adds a `:capture_gap` marker. Neither represents full evidence.
failed events retain a small error marker with validated HTTP status, provider
code and retryability when available and space permits, never error messages
or request bodies.
Snapshot eviction adds a `:capture_gap` marker. Neither represents full evidence.
A sink receives all bounded events; the snapshot is a bounded recent window.
"""

Expand Down Expand Up @@ -465,15 +468,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 length(state.events) > state.limits.max_events or
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 @@ -82,6 +82,109 @@ defmodule Imp.RunObservationTest do
assert terminal.kind == :run_cancelled
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