Skip to content
Merged
Show file tree
Hide file tree
Changes from 8 commits
Commits
Show all changes
47 commits
Select commit Hold shift + click to select a range
a9ae3b0
fix(stream-events): enforce timeout_secs after a silent backend
mkoushni Sep 9, 2026
a2db22e
Merge remote-tracking branch 'origin/main' into fix/938-stream-events…
mkoushni Sep 9, 2026
bebbc4e
fix(stream-events): keep completed streams persistable after idle close
mkoushni Sep 9, 2026
7020707
fix(tests): keep agentic-loop read_timeout unique and shrink chunked …
mkoushni Sep 9, 2026
3f77b40
Merge branch 'main' into fix/938-stream-events-idle-timeout
mkoushni Sep 10, 2026
3a009e6
Merge branch 'main' into fix/938-stream-events-idle-timeout
leseb Sep 10, 2026
f561a9b
Merge branch 'main' into fix/938-stream-events-idle-timeout
mkoushni Sep 10, 2026
691e206
Merge branch 'main' into fix/938-stream-events-idle-timeout
mkoushni Sep 10, 2026
1f1e511
Merge branch 'main' into fix/938-stream-events-idle-timeout
mkoushni Sep 10, 2026
d0068c3
fix(stream-events): skip idle-timeout errors after a terminal SSE event
mkoushni Sep 14, 2026
b2be8a0
Merge origin/main into fix/938-stream-events-idle-timeout
mkoushni Sep 14, 2026
7fb53b3
fix(stream-events): keep idle-timeout tests valid after IRR wrap
mkoushni Sep 14, 2026
62359ed
Merge origin/main into fix/938-stream-events-idle-timeout
mkoushni Sep 14, 2026
c1b7b10
Merge origin/main into fix/938-stream-events-idle-timeout
mkoushni Sep 14, 2026
a094ff8
fix(tests): allow logical SSE error after a committed stream
mkoushni Sep 14, 2026
3ad5cf8
Merge remote-tracking branch 'origin/main' into fix/938-stream-events…
mkoushni Sep 14, 2026
529ff44
fix(stream-events): apply timeout_secs after load balancing
mkoushni Sep 14, 2026
c4486fb
fix(docs): sync stream_events body-access table
mkoushni Sep 14, 2026
94c26bb
chore(deps): bump rustls to 0.23.45
mkoushni Sep 14, 2026
c83d70a
Merge remote-tracking branch 'origin/main' into fix/938-stream-events…
mkoushni Sep 15, 2026
a8fce97
fix(stream-events): recap leftover timeout on restored IRR peer
mkoushni Sep 15, 2026
7184838
test(openresponses): serialize live templates to avoid CPU starvation
mkoushni Sep 15, 2026
48689a8
docs(stream-events): leftover recap is copied onto the live body
mkoushni Sep 15, 2026
51180b1
Merge branch 'main' into fix/938-stream-events-idle-timeout
mkoushni Sep 15, 2026
196d738
Merge branch 'main' into fix/938-stream-events-idle-timeout
mkoushni Sep 16, 2026
c66094c
ci(integration): harden pinned Codex CLI download
mkoushni Sep 16, 2026
6eddb18
Merge branch 'main' into fix/938-stream-events-idle-timeout
mkoushni Sep 16, 2026
f14b2fa
fix(stream_events): start timeout at first chunk and recap live body
mkoushni Sep 16, 2026
05fde2f
fix(stream_events): pin Praxis 1172 and close timeout recap gaps
mkoushni Sep 16, 2026
c40ddf3
Merge branch 'main' into fix/938-stream-events-idle-timeout
mkoushni Sep 16, 2026
49ede10
Merge branch 'main' into fix/938-stream-events-idle-timeout
mkoushni Sep 16, 2026
891cbf9
Fix broken rustdoc link for cap_stream_read_timeout
mkoushni Sep 16, 2026
8031631
Regenerate openai_stream_events filter docs
mkoushni Sep 16, 2026
2e46338
Merge branch 'main' into fix/938-stream-events-idle-timeout
mkoushni Sep 16, 2026
4c6ee1f
fix(stream_events): cap absolute deadline on live body after first chunk
mkoushni Sep 16, 2026
73b6444
fix(ci): allow temporary mkoushni/praxis git pin in deny.toml
mkoushni Sep 16, 2026
3d19953
fix(deps): pin praxis core/filter via upstream git, not a personal fork
mkoushni Sep 16, 2026
2dd5299
fix(deps): git-pin all praxis crates to avoid duplicate core types
mkoushni Sep 16, 2026
382849c
merge: update stream-events branch with main
mkoushni Sep 17, 2026
7a03e3d
fix(deps): pin Praxis to v0.5.6
mkoushni Sep 17, 2026
67c2877
fix(tests): import json_post for stream events
mkoushni Sep 17, 2026
461efd5
fix(stream-events): adapt timeout tests to Praxis 0.5.6
mkoushni Sep 17, 2026
e8551ee
docs: regenerate filter documentation
mkoushni Sep 17, 2026
ac49537
fix(xtask): make filter doc discovery deterministic
mkoushni Sep 17, 2026
7f57277
Merge branch 'main' into fix/938-stream-events-idle-timeout
mkoushni Sep 17, 2026
747c9e5
fix(docs): use resolved Praxis source for filter docs
mkoushni Sep 17, 2026
8aeb9e9
Merge branch 'main' into fix/938-stream-events-idle-timeout
mkoushni Sep 17, 2026
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 apis/src/openai/responses/README.md
Original file line number Diff line number Diff line change
Expand Up @@ -42,7 +42,7 @@
| `openai_responses_proxy` | — | ReadWrite / StreamBuffer | — | — |
| `openai_responses_rehydrate` | — | ReadOnly / StreamBuffer | ✓ | ReadWrite / Stream |
| `openai_responses_validate` | — | ReadOnly / StreamBuffer | — | — |
| `openai_stream_events` | ✓ | | ✓ | ReadWrite / Stream |
| `openai_stream_events` | ✓ | ReadOnly / Stream | ✓ | ReadWrite / Stream |

Check warning on line 45 in apis/src/openai/responses/README.md

View workflow job for this annotation

GitHub Actions / Detect hidden unicode characters

Unicode Safety [non-ascii-identifier]

U+2713 <unnamed U+2713> -- Non-ASCII U+2713 <unnamed U+2713> in identifier '✓' (policy: ascii-only)

Check warning on line 45 in apis/src/openai/responses/README.md

View workflow job for this annotation

GitHub Actions / Detect hidden unicode characters

Unicode Safety [non-ascii-identifier]

U+2713 <unnamed U+2713> -- Non-ASCII U+2713 <unnamed U+2713> in identifier '✓' (policy: ascii-only)
| `openai_tool_parse` | ✓ | ReadOnly / StreamBuffer | — | — |
| `openai_web_search` | — | ReadOnly / StreamBuffer | — | ReadOnly / Stream |
| `responses_to_chat_completions` | — | ReadWrite / StreamBuffer | ✓ | ReadWrite / Stream |
7 changes: 7 additions & 0 deletions apis/src/openai/responses/stream_events/config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -34,6 +34,13 @@ pub(crate) struct StreamEventsConfig {
pub max_events: Option<usize>,

/// Maximum seconds from first chunk to stream completion.
///
/// Checked when SSE chunks or end-of-stream arrive. An idle backend
/// that sends nothing further does not invoke those callbacks, so this
/// budget is also applied as an upstream `read_timeout` when the filter
/// can see the selected peer, and must be paired with cluster
/// `read_timeout_ms` no larger than this value so a silent connection
/// is torn down without waiting for another chunk.
/// Default: 300 (5 minutes).
#[serde(default)]
pub timeout_secs: Option<u64>,
Expand Down
170 changes: 144 additions & 26 deletions apis/src/openai/responses/stream_events/mod.rs
Original file line number Diff line number Diff line change
Expand Up @@ -17,14 +17,15 @@ mod config;
use std::{
collections::{BTreeSet, hash_map::DefaultHasher},
hash::{Hash as _, Hasher as _},
sync::Arc,
time::{Duration, Instant},
};

use async_trait::async_trait;
use bytes::Bytes;
use praxis_filter::{
BodyAccess, BodyMode, FilterAction, FilterError, HttpFilter, HttpFilterContext, SubRequestResponseMode,
parse_filter_config,
BodyAccess, BodyMode, FilterAction, FilterError, HttpFilter, HttpFilterContext, StreamTerminationCause,
SubRequestResponseMode, parse_filter_config,
};
use serde_json::Value;
use tracing::{debug, trace, warn};
Expand Down Expand Up @@ -117,6 +118,10 @@ pub(super) struct StreamEventsState {
/// # timeout_secs: 300
/// # max_tool_call_argument_bytes: 1048576
/// ```
///
/// `timeout_secs` bounds elapsed time across chunks. Pair it with the
/// cluster's `read_timeout_ms` so a backend that goes silent after the
/// first event is still terminated.
pub struct OpenaiStreamEventsFilter {
/// Configuration for the SSE frame parser.
parser_config: SseParserConfig,
Expand Down Expand Up @@ -196,7 +201,10 @@ impl HttpFilter for OpenaiStreamEventsFilter {
}

fn request_body_access(&self) -> BodyAccess {
BodyAccess::None
// Observe the request-body phase so the timeout can be applied
// after load balancing has selected an upstream (when that phase
// still runs after `on_request`).
BodyAccess::ReadOnly
}

fn request_body_mode(&self) -> BodyMode {
Expand Down Expand Up @@ -229,11 +237,24 @@ impl HttpFilter for OpenaiStreamEventsFilter {
// failing an otherwise valid request. Strip `Accept-Encoding` whenever
// logical parsing is armed so a compliant backend returns plaintext SSE.
ctx.request_headers_to_remove.push(http::header::ACCEPT_ENCODING);
cap_upstream_read_timeout(ctx, self.parser_config.timeout);
}

Ok(FilterAction::Continue)
}

async fn on_request_body(
&self,
ctx: &mut HttpFilterContext<'_>,
_body: &mut Option<Bytes>,
end_of_stream: bool,
) -> Result<FilterAction, FilterError> {
if end_of_stream && Self::is_armed(ctx) {
cap_upstream_read_timeout(ctx, self.parser_config.timeout);
}
Ok(FilterAction::Continue)
}

async fn on_response(&self, ctx: &mut HttpFilterContext<'_>) -> Result<FilterAction, FilterError> {
if !Self::is_armed(ctx) {
return Ok(FilterAction::Continue);
Expand Down Expand Up @@ -273,6 +294,7 @@ impl HttpFilter for OpenaiStreamEventsFilter {
process_chunk(ctx, body);

if end_of_stream {
record_idle_transport_timeout(ctx);
validate_stream_end(ctx);
finalize_logical_stream(ctx, body);
}
Expand All @@ -281,6 +303,56 @@ impl HttpFilter for OpenaiStreamEventsFilter {
}
}

/// Cap the selected peer's per-read timeout at the stream budget.
///
/// Pingora and IRR streaming reads wake on this timer even when the
/// backend sends no further SSE bytes. A tighter cluster `read_timeout`
/// is left in place. No-op until load balancing has set `ctx.upstream`.
fn cap_upstream_read_timeout(ctx: &mut HttpFilterContext<'_>, timeout: Duration) {
let Some(upstream) = ctx.upstream.as_mut() else {
return;
};
let opts = Arc::make_mut(&mut upstream.connection);
opts.read_timeout = Some(opts.read_timeout.map_or(timeout, |existing| existing.min(timeout)));
}

/// Treat an IRR idle or deadline abort as a stream timeout.
///
/// Those failures arrive as end-of-stream with [`StreamTerminationCause`]
/// set, not as another SSE chunk, so [`check_timeout`] never ran while
/// the backend was silent. A backend that already sent a terminal event
/// can still trip this timer while closing the HTTP body; that is not a
/// stream error.
fn record_idle_transport_timeout(ctx: &mut HttpFilterContext<'_>) {
let Some(cause) = ctx.stream_termination().map(praxis_filter::StreamTermination::cause) else {
return;
};
if !matches!(
cause,
StreamTerminationCause::IdleTimeout | StreamTerminationCause::DeadlineExceeded
Comment thread
mkoushni marked this conversation as resolved.
Outdated
) {
return;
}
record_idle_timeout_error_if_incomplete(ctx);
ctx.mark_stream_termination_handled();
}

/// Set timeout error metadata only when the SSE parser never saw a terminal event.
fn record_idle_timeout_error_if_incomplete(ctx: &mut HttpFilterContext<'_>) {
let parser_complete = ctx
.get_filter_state::<StreamEventsState>()
.is_some_and(|state| state.completion_state != CompletionState::Open);
if parser_complete || ctx.get_metadata("responses.stream_error_code").is_some() {
return;
}
ctx.set_metadata("responses.stream_error_code", "server_error");
ctx.set_metadata(
"responses.stream_error_message",
"upstream Responses stream exceeded timeout",
);
ctx.set_metadata("responses.skip_persist", "true");
}

/// Parse SSE frames, accumulating state and optionally normalizing output.
fn process_chunk(ctx: &mut HttpFilterContext<'_>, body: &mut Option<Bytes>) {
let Some(bytes) = body.as_ref() else {
Expand Down Expand Up @@ -338,7 +410,11 @@ fn handle_parse_error(
ctx.set_metadata("responses.stream_error_code", "server_error");
ctx.set_metadata(
"responses.stream_error_message",
"upstream Responses stream could not be parsed",
if matches!(error, SseParseError::Timeout { .. }) {
"upstream Responses stream exceeded timeout"
} else {
"upstream Responses stream could not be parsed"
},
);
ctx.set_metadata("responses.skip_persist", "true");
*body = None;
Expand Down Expand Up @@ -1276,32 +1352,74 @@ fn mark_complete(state: &mut StreamEventsState, new_state: CompletionState, now:

/// Check that the SSE stream terminated with a terminal event.
fn validate_stream_end(ctx: &mut HttpFilterContext<'_>) {
let incomplete_logical_stream = ctx.get_filter_state::<StreamEventsState>().and_then(|state| {
let checked_at = state.completed_at.unwrap_or_else(Instant::now);
if let Err(e) = check_timeout(state, checked_at) {
warn!(error = %e, "stream did not terminate cleanly");
Some(state.logical_stream)
} else if state.completion_state == CompletionState::Open {
warn!("stream did not terminate cleanly: missing terminal event");
Some(state.logical_stream)
} else {
None
}
});
if let Some(logical_stream) = incomplete_logical_stream {
ctx.set_metadata("responses.stream_incomplete", "true".to_owned());
if logical_stream && ctx.get_metadata("responses.stream_error_code").is_none() {
ctx.set_metadata("responses.stream_error_code", "server_error");
ctx.set_metadata(
"responses.stream_error_message",
"upstream Responses stream did not terminate cleanly",
);
ctx.set_metadata("responses.skip_persist", "true");
}
match stream_end_kind(ctx) {
StreamEndKind::Complete => {},
StreamEndKind::Incomplete {
logical_stream,
timed_out,
} => {
record_incomplete_stream(ctx, logical_stream, timed_out);
},
}
debug!("stream_events processing complete");
}

/// Classify how the current parser state ended.
fn stream_end_kind(ctx: &HttpFilterContext<'_>) -> StreamEndKind {
let Some(state) = ctx.get_filter_state::<StreamEventsState>() else {
return StreamEndKind::Complete;
};
let checked_at = state.completed_at.unwrap_or_else(Instant::now);
match check_timeout(state, checked_at) {
Err(error) => {
warn!(%error, "stream did not terminate cleanly");
StreamEndKind::Incomplete {
logical_stream: state.logical_stream,
timed_out: true,
}
},
Ok(()) if state.completion_state == CompletionState::Open => {
warn!("stream did not terminate cleanly: missing terminal event");
StreamEndKind::Incomplete {
logical_stream: state.logical_stream,
timed_out: false,
}
},
Ok(()) => StreamEndKind::Complete,
}
}

/// Publish incomplete-stream metadata for persistence and logical-stream errors.
fn record_incomplete_stream(ctx: &mut HttpFilterContext<'_>, logical_stream: bool, timed_out: bool) {
ctx.set_metadata("responses.stream_incomplete", "true".to_owned());
if !logical_stream || ctx.get_metadata("responses.stream_error_code").is_some() {
return;
}
ctx.set_metadata("responses.stream_error_code", "server_error");
ctx.set_metadata(
"responses.stream_error_message",
if timed_out {
"upstream Responses stream exceeded timeout"
} else {
"upstream Responses stream did not terminate cleanly"
},
);
ctx.set_metadata("responses.skip_persist", "true");
}

/// How an SSE stream ended from the parser's point of view.
enum StreamEndKind {
/// A terminal lifecycle or error event was observed in time.
Complete,
/// The stream ended without a clean terminal event.
Incomplete {
/// Whether this parser is normalizing an IRR logical stream.
logical_stream: bool,
/// Whether the wall-clock budget was exceeded.
timed_out: bool,
},
}

/// Whether the response is a successful `text/event-stream` response.
fn is_success_sse_response(ctx: &HttpFilterContext<'_>) -> bool {
let Some(resp) = ctx.response_header.as_ref() else {
Expand Down
106 changes: 106 additions & 0 deletions apis/src/openai/responses/stream_events/tests.rs
Original file line number Diff line number Diff line change
Expand Up @@ -3371,3 +3371,109 @@ async fn on_response_preserves_content_length_when_not_armed() {
"Content-Length should be preserved when filter is not armed"
);
}

#[tokio::test]
async fn on_request_body_caps_upstream_read_timeout_at_stream_budget() {
use std::{sync::Arc, time::Duration};

use praxis_core::connectivity::{ConnectionOptions, Upstream};

let yaml: serde_yaml::Value = serde_yaml::from_str("timeout_secs: 1").unwrap();
let filter = OpenaiStreamEventsFilter::from_config(&yaml).unwrap();
let req = make_request(http::Method::POST, "/v1/responses");
let mut ctx = make_filter_context(Box::leak(Box::new(req)));
ctx.set_metadata("openai_responses_format.format", "openai_responses".to_owned());
ctx.set_metadata("openai_responses_format.stream", "true".to_owned());
ctx.current_filter_id = Some(0);
ctx.upstream = Some(Upstream {
address: Arc::from("127.0.0.1:9"),
authority: None,
tls: None,
connection: Arc::new(ConnectionOptions::default()),
});

filter.on_request(&mut ctx).await.unwrap();
filter.on_request_body(&mut ctx, &mut None, true).await.unwrap();

assert_eq!(
ctx.upstream.as_ref().and_then(|u| u.connection.read_timeout),
Some(Duration::from_secs(1)),
"armed stream_events must cap the selected peer read timeout so idle backends wake"
);
}

#[tokio::test]
async fn on_request_body_keeps_a_tighter_cluster_read_timeout() {
use std::{sync::Arc, time::Duration};

use praxis_core::connectivity::{ConnectionOptions, Upstream};

let yaml: serde_yaml::Value = serde_yaml::from_str("timeout_secs: 300").unwrap();
let filter = OpenaiStreamEventsFilter::from_config(&yaml).unwrap();
let req = make_request(http::Method::POST, "/v1/responses");
let mut ctx = make_filter_context(Box::leak(Box::new(req)));
ctx.set_metadata("openai_responses_format.format", "openai_responses".to_owned());
ctx.set_metadata("openai_responses_format.stream", "true".to_owned());
ctx.current_filter_id = Some(0);
ctx.upstream = Some(Upstream {
address: Arc::from("127.0.0.1:9"),
authority: None,
tls: None,
connection: Arc::new(ConnectionOptions {
read_timeout: Some(Duration::from_millis(250)),
..ConnectionOptions::default()
}),
});

filter.on_request(&mut ctx).await.unwrap();

assert_eq!(
ctx.upstream.as_ref().and_then(|u| u.connection.read_timeout),
Some(Duration::from_millis(250)),
"a tighter cluster read timeout must not be relaxed to timeout_secs"
);
}

#[tokio::test]
async fn idle_timeout_does_not_fail_a_completed_sse_stream() {
let (filter, mut ctx) = make_armed_context();
filter.on_request(&mut ctx).await.unwrap();

let completed =
json!({"id": "resp_1", "status": "completed", "model": "m", "created_at": 0, "output": [], "usage": {}});
let mut body = Some(make_sse_chunk("response.completed", &completed));
filter.on_response_body(&mut ctx, &mut body, false).unwrap();

super::record_idle_timeout_error_if_incomplete(&mut ctx);

assert!(
ctx.get_metadata("responses.stream_error_code").is_none(),
"a terminal lifecycle event must not be rewritten as a transport timeout"
);
assert!(
ctx.get_metadata("responses.skip_persist").is_none(),
"a completed SSE stream must remain persistable when the HTTP body closes slowly"
);
}

#[tokio::test]
async fn idle_timeout_fails_an_open_sse_stream() {
let (filter, mut ctx) = make_armed_context();
filter.on_request(&mut ctx).await.unwrap();

let mut body = Some(make_sse_chunk("response.output_text.delta", &json!({"text": "hi"})));
filter.on_response_body(&mut ctx, &mut body, false).unwrap();

super::record_idle_timeout_error_if_incomplete(&mut ctx);

assert_eq!(
ctx.get_metadata("responses.stream_error_code"),
Some("server_error"),
"an idle abort before a terminal event is a stream timeout"
);
assert_eq!(
ctx.get_metadata("responses.skip_persist"),
Some("true"),
"an incomplete idle abort must not be persisted"
);
}
4 changes: 4 additions & 0 deletions apis/src/openai/sse/config.rs
Original file line number Diff line number Diff line change
Expand Up @@ -27,6 +27,10 @@ pub(crate) struct SseParserConfig {
pub max_events: usize,

/// Maximum wall-clock time from first chunk to stream completion.
///
/// Enforced when chunks or end-of-stream arrive. Idle gaps with no
/// body traffic are enforced by the upstream read timeout (see
/// `openai_stream_events.timeout_secs`).
pub timeout: Duration,
}

Expand Down
2 changes: 1 addition & 1 deletion docs/filters/openai_stream_events.md
Original file line number Diff line number Diff line change
Expand Up @@ -16,7 +16,7 @@ All fields are optional; omitted values fall back to [`SseParserConfig`] default
| `logical_stream` | bool | no | Treat successive IRR inference streams as one logical Responses stream. Per-iteration lifecycle events are normalized and only the final terminal event is exposed downstream. |
| `max_buffer_bytes` | integer | no | Maximum bytes buffered for incomplete SSE lines/data across chunk boundaries. Default: 10 MiB. |
| `max_events` | integer | no | Maximum number of SSE events before the parser errors. Default: 100,000. |
| `timeout_secs` | integer | no | Maximum seconds from first chunk to stream completion. Default: 300 (5 minutes). |
| `timeout_secs` | integer | no | Maximum seconds from first chunk to stream completion. Checked when SSE chunks or end-of-stream arrive. An idle backend that sends nothing further does not invoke those callbacks, so this budget is also applied as an upstream `read_timeout` when the filter can see the selected peer, and must be paired with cluster `read_timeout_ms` no larger than this value so a silent connection is torn down without waiting for another chunk. Default: 300 (5 minutes). |
| `max_tool_call_argument_bytes` | integer | no | Maximum bytes accepted per function-call argument string from `function_call_arguments.delta` or `function_call_arguments.done` events. Default: 1 MiB. |

## Example
Expand Down
Loading
Loading