Merge remote-tracking branch 'origin/pr-483' into merge-pr-483

# Conflicts:
#	apps/aether-gateway/src/ai_serving/planner/standard/openai/chat/decision/payload.rs
#	apps/aether-gateway/src/ai_serving/planner/standard/openai/chat/decision/request.rs
#	apps/aether-gateway/src/ai_serving/planner/standard/openai/responses/decision/payload.rs
#	apps/aether-gateway/src/ai_serving/planner/standard/openai/responses/decision/request.rs
#	apps/aether-gateway/src/execution_runtime/chatgpt_web_image.rs
This commit is contained in:
fawney19
2026-05-19 02:23:14 +08:00
53 changed files with 8081 additions and 305 deletions
@@ -1044,6 +1044,12 @@ fn build_sse_body_stream(
}
}
fn stream_chunk_contains_sse_done(chunk: &[u8]) -> bool {
std::str::from_utf8(chunk)
.ok()
.is_some_and(|text| text.lines().any(|line| line.trim() == "data: [DONE]"))
}
async fn next_stream_frame<R>(
buffered_frames: &mut VecDeque<StreamFrame>,
lines: &mut FramedRead<R, LinesCodec>,
@@ -1088,6 +1094,33 @@ fn should_refresh_stream_usage_telemetry(
|| (next_elapsed.is_some() && next_elapsed != previous_elapsed)
}
fn build_terminal_stream_telemetry(
stream_started_at: Instant,
telemetry: Option<&ExecutionTelemetry>,
usage_stream_telemetry: Option<&ExecutionTelemetry>,
upstream_bytes: u64,
) -> ExecutionTelemetry {
let current_elapsed_ms = stream_started_at
.elapsed()
.as_millis()
.min(u128::from(u64::MAX)) as u64;
let ttfb_ms = telemetry
.and_then(|telemetry| telemetry.ttfb_ms)
.or_else(|| usage_stream_telemetry.and_then(|telemetry| telemetry.ttfb_ms));
let prior_elapsed_ms = telemetry
.and_then(|telemetry| telemetry.elapsed_ms)
.or_else(|| usage_stream_telemetry.and_then(|telemetry| telemetry.elapsed_ms))
.unwrap_or(0);
let elapsed_ms = current_elapsed_ms
.max(prior_elapsed_ms)
.max(ttfb_ms.unwrap_or(0));
ExecutionTelemetry {
ttfb_ms,
elapsed_ms: Some(elapsed_ms),
upstream_bytes: Some(upstream_bytes),
}
}
fn should_skip_direct_finalize_prefetch(
direct_stream_finalize_kind: Option<&str>,
content_type: Option<&str>,
@@ -2052,6 +2085,8 @@ async fn execute_stream_from_frame_stream(
max_stream_body_buffer_bytes,
&mut client_body_truncated,
);
let mut client_visible_stream_completed =
stream_chunk_contains_sse_done(&prefetched_body_for_report);
let mut usage_stream_telemetry: Option<ExecutionTelemetry> = initial_telemetry.clone();
let mut telemetry: Option<ExecutionTelemetry> = initial_telemetry;
let reached_eof = initial_reached_eof;
@@ -2478,6 +2513,8 @@ async fn execute_stream_from_frame_stream(
);
let rewritten_chunk_len =
u64::try_from(rewritten_chunk.len()).unwrap_or(u64::MAX);
let chunk_completed_stream =
stream_chunk_contains_sse_done(&rewritten_chunk);
if tx.send(Ok(Bytes::from(rewritten_chunk))).await.is_err() {
warn!(
event_name = "stream_execution_downstream_disconnected",
@@ -2490,6 +2527,7 @@ async fn execute_stream_from_frame_stream(
downstream_dropped = true;
break;
} else {
client_visible_stream_completed |= chunk_completed_stream;
client_stream_bytes.fetch_add(rewritten_chunk_len, Ordering::Relaxed);
last_client_chunk_elapsed_ms.store(
stream_started_at_for_report
@@ -2604,6 +2642,8 @@ async fn execute_stream_from_frame_stream(
);
let rewritten_chunk_len =
u64::try_from(rewritten_chunk.len()).unwrap_or(u64::MAX);
let chunk_completed_stream =
stream_chunk_contains_sse_done(&rewritten_chunk);
if tx.send(Ok(Bytes::from(rewritten_chunk))).await.is_err() {
warn!(
event_name = "stream_execution_downstream_flush_disconnected",
@@ -2615,6 +2655,7 @@ async fn execute_stream_from_frame_stream(
);
downstream_dropped = true;
} else {
client_visible_stream_completed |= chunk_completed_stream;
client_stream_bytes
.fetch_add(rewritten_chunk_len, Ordering::Relaxed);
last_client_chunk_elapsed_ms.store(
@@ -2661,6 +2702,8 @@ async fn execute_stream_from_frame_stream(
);
let flushed_chunk_len =
u64::try_from(flushed_chunk.len()).unwrap_or(u64::MAX);
let chunk_completed_stream =
stream_chunk_contains_sse_done(&flushed_chunk);
if tx.send(Ok(Bytes::from(flushed_chunk))).await.is_err() {
warn!(
event_name = "stream_execution_downstream_rewrite_flush_disconnected",
@@ -2672,6 +2715,7 @@ async fn execute_stream_from_frame_stream(
);
downstream_dropped = true;
} else {
client_visible_stream_completed |= chunk_completed_stream;
client_stream_bytes.fetch_add(flushed_chunk_len, Ordering::Relaxed);
last_client_chunk_elapsed_ms.store(
stream_started_at_for_report
@@ -2781,6 +2825,18 @@ async fn execute_stream_from_frame_stream(
),
);
if downstream_dropped && client_visible_stream_completed && terminal_failure.is_none() {
debug!(
event_name = "execution_runtime_stream_downstream_closed_after_done",
log_type = "debug",
trace_id = %trace_id_owned,
request_id = %request_id_for_report_log,
candidate_id = ?candidate_id_for_report.as_deref(),
"gateway treats downstream close after client-visible SSE DONE as completed"
);
downstream_dropped = false;
}
if downstream_dropped {
debug!(
event_name = "execution_runtime_stream_report_skipped",
@@ -2791,6 +2847,12 @@ async fn execute_stream_from_frame_stream(
trace_id = %trace_id_owned,
"gateway skipped stream report because downstream disconnected before completion"
);
let terminal_telemetry = Some(build_terminal_stream_telemetry(
stream_started_at_for_report,
telemetry.as_ref(),
usage_stream_telemetry.as_ref(),
provider_stream_bytes.load(Ordering::Relaxed),
));
let usage_payload = build_stream_usage_payload(
trace_id_owned,
report_kind_owned.unwrap_or_default(),
@@ -2802,7 +2864,7 @@ async fn execute_stream_from_frame_stream(
&buffered_body,
client_body_truncated,
stream_terminal_summary,
telemetry,
terminal_telemetry,
);
record_stream_terminal_usage(
&state_for_report,
@@ -2834,6 +2896,12 @@ async fn execute_stream_from_frame_stream(
if let Some(failure) = terminal_failure {
record_manual_proxy_stream_error(&state_for_report, &plan_for_report).await;
let terminal_telemetry = Some(build_terminal_stream_telemetry(
stream_started_at_for_report,
telemetry.as_ref(),
usage_stream_telemetry.as_ref(),
provider_stream_bytes.load(Ordering::Relaxed),
));
submit_midstream_stream_failure(
&state_for_report,
&trace_id_owned,
@@ -2841,7 +2909,7 @@ async fn execute_stream_from_frame_stream(
direct_stream_finalize_kind_owned.as_deref(),
report_context_owned,
headers_for_report,
telemetry,
terminal_telemetry,
&provider_buffered_body,
candidate_started_unix_secs_for_report,
failure,
@@ -2851,6 +2919,12 @@ async fn execute_stream_from_frame_stream(
}
let should_submit_report = report_kind_owned.is_some();
let terminal_telemetry = Some(build_terminal_stream_telemetry(
stream_started_at_for_report,
telemetry.as_ref(),
usage_stream_telemetry.as_ref(),
provider_stream_bytes.load(Ordering::Relaxed),
));
let usage_payload = build_stream_usage_payload(
trace_id_owned.clone(),
report_kind_owned.unwrap_or_default(),
@@ -2862,7 +2936,7 @@ async fn execute_stream_from_frame_stream(
&buffered_body,
client_body_truncated,
stream_terminal_summary,
telemetry,
terminal_telemetry,
);
apply_local_execution_effect(
&state_for_report,
@@ -3365,6 +3439,7 @@ mod tests {
.await
.expect("first business chunk should arrive");
assert_eq!(first.as_ref(), b"data: {\"id\":\"first\"}\n\n");
tokio::time::sleep(Duration::from_millis(30)).await;
drop(body_stream);
tokio::time::timeout(Duration::from_secs(1), frame_stream_dropped.notified())
@@ -3392,6 +3467,172 @@ mod tests {
candidates[0].error_type.as_deref(),
Some("downstream_disconnect")
);
let stored_usage = tokio::time::timeout(Duration::from_secs(1), async {
loop {
let usage = usage_repository
.find_by_request_id("req-client-drop-cancels-upstream")
.await
.expect("usage should read");
if usage
.as_ref()
.is_some_and(|usage| usage.status == "cancelled")
{
break usage.expect("cancelled usage should exist");
}
tokio::time::sleep(Duration::from_millis(10)).await;
}
})
.await
.expect("usage should be marked cancelled");
assert_eq!(stored_usage.billing_status, "pending");
assert_eq!(stored_usage.status_code, Some(499));
let first_byte_time_ms = stored_usage
.first_byte_time_ms
.expect("cancelled stream should retain first byte time");
let response_time_ms = stored_usage
.response_time_ms
.expect("cancelled stream should record terminal duration");
assert!(
response_time_ms > first_byte_time_ms,
"terminal duration should include time after the first byte"
);
}
#[tokio::test]
async fn image_stream_downstream_close_after_done_is_recorded_success() {
let usage_repository = Arc::new(InMemoryUsageReadRepository::default());
let request_candidate_repository = Arc::new(InMemoryRequestCandidateRepository::default());
let state = AppState::new()
.expect("app state should build")
.with_data_state_for_tests(
crate::data::GatewayDataState::with_request_candidate_and_usage_repository_for_tests(
Arc::clone(&request_candidate_repository),
Arc::clone(&usage_repository),
),
)
.with_usage_runtime_for_tests(UsageRuntimeConfig {
enabled: true,
..UsageRuntimeConfig::default()
});
let plan = ExecutionPlan {
request_id: "req-image-done-close-success".into(),
candidate_id: Some("cand-image-done-close-success".into()),
provider_name: Some("openai".into()),
provider_id: "prov-1".into(),
endpoint_id: "ep-1".into(),
key_id: "key-1".into(),
method: "POST".into(),
url: "https://example.com/v1/images/generations".into(),
headers: BTreeMap::from([("accept".into(), "text/event-stream".into())]),
content_type: Some("application/json".into()),
content_encoding: None,
body: RequestBody::from_json(json!({
"model": "gpt-image-2",
"prompt": "draw a small image",
"stream": true
})),
stream: true,
client_api_format: "openai:chat".into(),
provider_api_format: "openai:image".into(),
model_name: Some("gpt-image-2".into()),
proxy: None,
transport_profile: None,
timeouts: None,
};
let frame_stream = stream! {
yield Ok::<Bytes, std::io::Error>(Bytes::from_static(
b"{\"type\":\"headers\",\"payload\":{\"kind\":\"headers\",\"status_code\":200,\"headers\":{\"content-type\":\"text/event-stream\"}}}\n",
));
yield Ok::<Bytes, std::io::Error>(Bytes::from_static(
b"{\"type\":\"data\",\"payload\":{\"kind\":\"data\",\"text\":\"event: response.output_item.done\\ndata: {\\\"type\\\":\\\"response.output_item.done\\\",\\\"output_index\\\":0,\\\"item\\\":{\\\"id\\\":\\\"ig_1\\\",\\\"type\\\":\\\"image_generation_call\\\",\\\"result\\\":\\\"aGVsbG8=\\\"}}\\n\\nevent: response.completed\\ndata: {\\\"type\\\":\\\"response.completed\\\",\\\"response\\\":{\\\"id\\\":\\\"resp_1\\\",\\\"model\\\":\\\"gpt-image-2\\\",\\\"status\\\":\\\"completed\\\",\\\"usage\\\":null}}\\n\\n\"}}\n",
));
std::future::pending::<()>().await;
}
.boxed();
let response = execute_stream_from_frame_stream(
&state,
plan,
"trace-image-done-close-success",
&test_decision(),
"openai_chat_stream",
Some("openai_chat_stream_success".to_string()),
Some(json!({
"request_id": "req-image-done-close-success",
"candidate_id": "cand-image-done-close-success",
"candidate_index": 0,
"retry_index": 0,
"provider_api_format": "openai:image",
"client_api_format": "openai:chat",
"image_request": {
"size": "1024x1024",
"quality": "medium"
}
})),
crate::clock::current_unix_ms(),
Instant::now(),
frame_stream,
None,
)
.await
.expect("execution should succeed")
.expect("execution should return a client response");
let mut body_stream = response.into_body().into_data_stream();
let mut body = Vec::new();
tokio::time::timeout(Duration::from_secs(1), async {
while !String::from_utf8_lossy(&body).contains("data: [DONE]") {
let chunk = body_stream
.next()
.await
.expect("body should yield until done")
.expect("chunk should be ok");
body.extend_from_slice(&chunk);
}
})
.await
.expect("final DONE should arrive");
drop(body_stream);
let candidates = tokio::time::timeout(Duration::from_secs(1), async {
loop {
let candidates = request_candidate_repository
.list_by_request_id("req-image-done-close-success")
.await
.expect("request candidates should read");
if candidates
.first()
.is_some_and(|candidate| candidate.status == RequestCandidateStatus::Success)
{
break candidates;
}
tokio::time::sleep(Duration::from_millis(10)).await;
}
})
.await
.expect("candidate should be marked success");
assert_eq!(candidates[0].status_code, Some(200));
let stored_usage = tokio::time::timeout(Duration::from_secs(1), async {
loop {
let usage = usage_repository
.find_by_request_id("req-image-done-close-success")
.await
.expect("usage should read");
if usage
.as_ref()
.is_some_and(|usage| usage.status == "completed")
{
break usage.expect("completed usage should exist");
}
tokio::time::sleep(Duration::from_millis(10)).await;
}
})
.await
.expect("usage should be marked completed");
assert_eq!(stored_usage.status_code, Some(200));
assert!(stored_usage.total_tokens > 0);
}
#[tokio::test]