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
2 changes: 1 addition & 1 deletion Cargo.lock

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

2 changes: 1 addition & 1 deletion crates/agentkit-provider-openai/Cargo.toml
Original file line number Diff line number Diff line change
Expand Up @@ -7,7 +7,7 @@ edition.workspace = true
license.workspace = true
repository.workspace = true
rust-version.workspace = true
version = "0.10.10"
version = "0.10.11"

[dependencies]
tokio-tungstenite = { version = "=0.29.0", default-features = false, features = ["handshake"] }
Expand Down
54 changes: 42 additions & 12 deletions crates/agentkit-provider-openai/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -120,9 +120,10 @@ WebSocket retries are intentionally more conservative than HTTP retries:
- Handshake status failures use the existing bounded retry policy and observer.
- A wrapped HTTP error before `response.created` (and before visible output)
may reconnect and retry only on a fresh socket. On reused sockets, errors lack
reliable request correlation and may belong to a previous turn, so they never
trigger automatic replay or authentication refresh. An accepted response is
never automatically replayed.
reliable request correlation and may belong to a previous turn, so generic
errors never trigger automatic replay or authentication refresh. The narrowly
scoped missing-continuation recovery is described below. An accepted response
is never automatically replayed.
- An interrupted send, socket EOF, receive failure, or timeout after sending is
**not replayed**: the server may already have accepted the request.
- Visible WebSocket output is never superseded/replayed, even when the consumer
Expand All @@ -140,15 +141,44 @@ have a 30-second ceiling; configured attempt, idle, logical retry-budget, and
cancellation bounds also apply. As with HTTP, configure `with_resilience` to set
stream idle and whole-turn deadlines.

### Current limitations

The **full credential-bound transcript is always authoritative and always sent**.
This release deliberately does not send `previous_response_id` or maintain an
incremental response-ID cache: reconstructing an exact prefix from normalized
text/tool/reasoning/image output without losing provider fields requires a
separate lossless compatibility proof. Reuse therefore saves connection setup,
not request transcript bytes. Compaction, changed inputs, and reconnection do not
risk a stale server-side prefix.
The optional [live continuation suite](../../docs/live-responses-websocket.md)
verifies incremental wire requests and server acceptance using environment credentials.
It is ignored by default and makes billable requests when explicitly enabled.

### Incremental continuation and limitations

The **full credential-bound transcript remains authoritative**. After a successful
`response.completed`, the live connection retains a bounded, memory-only checkpoint:
the response ID, full encoded request, raw completed output items, and their normal
transcript replay encoding. When all non-input request fields match and the next
input begins with the previous input plus that replay output, the next
`response.create` sends `previous_response_id` and only the new input suffix.
This includes tool results and later user messages, including rounds containing
encrypted reasoning or generated images. No continuation state is persisted.

Compatibility uses the same credential-bound encoder as an ordinary full request.
In particular, function arguments are serialized from parsed JSON, message output
metadata such as status and annotations is not replayed, and reasoning replays its
ID and encrypted content with an empty summary. These are the adapter's transcript
semantics, not byte equality against the server's raw output. Changes to retained
text, call IDs/arguments, protected reasoning or image data fail the prefix proof.
Unreadable/unrepresentable or oversized checkpoints disable the optimization.
Changes to model, tools, instructions, reasoning settings, or compacted/edited
history send the full input without an ID. Reconnection, authentication changes,
cancellation and failures discard the connection checkpoint. Checkpoint request
and output buffers are bounded by the configured request/item limits and zeroized
on drop; the socket also bounds completed response-ID correlation history.

An explicit `previous_response_not_found` rejection before acceptance or visible
output may reconnect and retry once with the full request, **within the configured
resilience retry budget**. Recovery requires the documented 400
`invalid_request_error` message naming the exact predecessor sent by this request
(`Previous response with id '<id>' not found.`). Code-only errors, older IDs, changed
message formats, and malformed error events fail closed: they cannot safely
correlate a rejection on a reused socket. Set `with_resilience` with `max_retries >= 1` to allow
this recovery. The full retry has no previous ID, so this special recovery cannot
repeat. Accepted responses, visible output, and ambiguous sends/disconnects are
never replayed. Recovery uses the existing retry observations and deadline budget.

WebSocket upgrades use a dedicated reqwest HTTP/1 client with redirects and
implicit HTTP retries disabled, using the existing reqwest TLS stack. A custom
Expand Down
67 changes: 62 additions & 5 deletions crates/agentkit-provider-openai/src/responses.rs
Original file line number Diff line number Diff line change
Expand Up @@ -739,10 +739,17 @@ impl OpenAIResponsesTurn {
} else {
let attempt = self.attempt.as_mut().expect("live attempt is open");
attempt.truncated.observe(&chunk);
match &attempt.body {
match &mut attempt.body {
LiveBody::Http(_) => attempt.decoder.push(&chunk),
LiveBody::WebSocket(_) => {
LiveBody::WebSocket(connection) => {
let result = attempt.decoder.push_json(&chunk);
if result.is_ok() {
connection.observe(
&chunk,
&attempt.decoder.state,
&self.context,
);
}
if result.is_ok() && attempt.decoder.state.terminal {
attempt.eof = true;
}
Expand All @@ -761,7 +768,7 @@ impl OpenAIResponsesTurn {
}
}
};
if let Err(failure) = result {
if let Err(mut failure) = result {
self.context.tracker.note_failure(&failure.error);
let websocket = self
.attempt
Expand All @@ -770,11 +777,23 @@ impl OpenAIResponsesTurn {
// A sent WebSocket request has no idempotency guarantee. Only explicit
// provider rejections before visible output on a fresh socket
// can be retried; reused sockets can deliver uncorrelated stale errors.
let missing_previous = self.attempt.as_ref().is_some_and(|attempt| {
matches!(&attempt.body, LiveBody::WebSocket(connection)
if connection.missing_previous(&attempt.decoder))
});
// A precise missing-continuation rejection is recoverable by dropping
// this socket and using the authoritative full request. The new socket
// has no checkpoint, so this exception cannot repeat. It shares the
// existing retry budget, observer and cancellation-safe reopen path.
if missing_previous {
failure.retryable = true;
}
if websocket
&& (self.attempt_output_emitted
|| !self.attempt.as_ref().is_some_and(|attempt| {
attempt.decoder.request_rejected
&& matches!(&attempt.body, LiveBody::WebSocket(connection)
missing_previous
|| attempt.decoder.request_rejected
&& matches!(&attempt.body, LiveBody::WebSocket(connection)
if connection.can_retry_rejection())
}))
{
Expand Down Expand Up @@ -2138,6 +2157,7 @@ struct ResponsesSseDecoder {
buffer_start: usize,
received: usize,
request_rejected: bool,
previous_response_missing: Option<String>,
max_attempt_bytes: usize,
state: ResponsesState,
}
Expand Down Expand Up @@ -2170,6 +2190,7 @@ impl ResponsesSseDecoder {
buffer_start: 0,
received: 0,
request_rejected: false,
previous_response_missing: None,
max_attempt_bytes: limits.max_attempt_bytes,
state: ResponsesState::new(
model,
Expand Down Expand Up @@ -2328,10 +2349,46 @@ impl ResponsesSseDecoder {
zeroize_encrypted_content(&mut value);
return Err(protocol_failure("Responses SSE event name/type mismatch"));
}
self.previous_response_missing = None;
let mut result = self.state.consume(&kind, &value);
if kind == "error"
&& let Err(failure) = &mut result
&& failure
.error
.provider_failure()
.is_some_and(|failure| failure.reason == ProviderFailureReason::ResponseFailed)
{
// Only a validated rejection may authorize recovery, never a sequence
// or lifecycle protocol failure containing an error-shaped payload.
let error = value.get("error").unwrap_or(&value);
if !self.state.created
&& value
.get("status")
.or_else(|| value.get("status_code"))
.and_then(Value::as_u64)
== Some(400)
&& error.get("type").and_then(Value::as_str) == Some("invalid_request_error")
&& error
.get("param")
.is_none_or(|param| param == "previous_response_id")
&& error
.get("code")
.or_else(|| value.get("code"))
.and_then(Value::as_str)
== Some("previous_response_not_found")
{
// https://developers.openai.com/api/docs/guides/websocket-mode#errors-to-handle
// The code alone is uncorrelated on a reused socket. Extract only
// the bounded ID from the provider's precise rejection message;
// never retain or expose arbitrary provider message text.
self.previous_response_missing = error
.get("message")
.and_then(Value::as_str)
.and_then(|message| message.strip_prefix("Previous response with id '"))
.and_then(|message| message.strip_suffix("' not found."))
.filter(|id| !id.is_empty() && id.len() <= 512)
.map(str::to_owned);
}
if let Some(status) = value
.get("status")
.or_else(|| value.get("status_code"))
Expand Down
Loading
Loading