From e58570d79d4fe47c087c97502640f4f788b788a3 Mon Sep 17 00:00:00 2001
From: elky
Date: Tue, 8 Sep 2026 23:11:37 +0800
Subject: [PATCH] feat(routing): make client disconnect behavior
strategy-scoped
---
Cargo.lock | 1 +
apps/aether-gateway/Cargo.toml | 1 +
.../src/execution_runtime/stream/execution.rs | 324 +++++++++--------
.../src/executor/orchestration.rs | 83 ++++-
apps/aether-gateway/src/handlers/proxy/mod.rs | 2 +-
.../proxy/websocket/responses/connection.rs | 40 ++-
.../proxy/websocket/responses/turn.rs | 7 +
apps/aether-gateway/src/lib.rs | 1 +
apps/aether-gateway/src/request_lifecycle.rs | 336 ++++++++++++++++++
apps/aether-gateway/src/routing/resolver.rs | 14 +-
.../src/tests/ai_execute/lifecycle.rs | 26 +-
crates/aether-billing/src/event_enrichment.rs | 186 ++++++----
.../src/repository/usage/metadata_policy.rs | 9 +
crates/aether-routing-core/src/model.rs | 35 ++
.../tests/responses_websocket_e2e.rs | 140 +++++++-
crates/aether-usage/runtime/src/record.rs | 32 ++
crates/aether-usage/runtime/src/settlement.rs | 74 +++-
docs/operations/client-disconnect-policy.md | 39 ++
.../routing/__tests__/routingPolicy.spec.ts | 16 +-
.../features/routing/utils/routingPolicy.ts | 3 +
frontend/src/i18n/legacy-admin-messages.ts | 2 +
frontend/src/views/admin/RoutingProfiles.vue | 25 +-
22 files changed, 1156 insertions(+), 240 deletions(-)
create mode 100644 apps/aether-gateway/src/request_lifecycle.rs
create mode 100644 docs/operations/client-disconnect-policy.md
diff --git a/Cargo.lock b/Cargo.lock
index 3823839e9..a2171f728 100644
--- a/Cargo.lock
+++ b/Cargo.lock
@@ -305,6 +305,7 @@ dependencies = [
"futures-util",
"hmac",
"http",
+ "http-body",
"http-body-util",
"hyper",
"hyper-util",
diff --git a/apps/aether-gateway/Cargo.toml b/apps/aether-gateway/Cargo.toml
index f664afcff..8337ac465 100644
--- a/apps/aether-gateway/Cargo.toml
+++ b/apps/aether-gateway/Cargo.toml
@@ -62,6 +62,7 @@ flate2.workspace = true
futures-util.workspace = true
hmac.workspace = true
http.workspace = true
+http-body = "1"
http-body-util = "0.1"
hyper = { version = "1", features = ["client", "server", "http1", "http2"] }
hyper-util = { version = "0.1", features = ["client-legacy", "client-pool", "server-auto", "service", "tokio"] }
diff --git a/apps/aether-gateway/src/execution_runtime/stream/execution.rs b/apps/aether-gateway/src/execution_runtime/stream/execution.rs
index 97b9d0f91..da0d04f08 100644
--- a/apps/aether-gateway/src/execution_runtime/stream/execution.rs
+++ b/apps/aether-gateway/src/execution_runtime/stream/execution.rs
@@ -13607,10 +13607,12 @@ mod tests {
}
#[tokio::test]
- async fn execute_stream_from_frame_stream_cancels_upstream_when_client_drops_body() {
- let usage_repository = Arc::new(InMemoryUsageReadRepository::default());
- let request_candidate_repository = Arc::new(InMemoryRequestCandidateRepository::default());
- let state = AppState::new()
+ async fn execute_stream_from_frame_stream_honors_client_disconnect_policy() {
+ for cancel_on_client_disconnect in [false, true] {
+ 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(
@@ -13622,39 +13624,39 @@ mod tests {
enabled: true,
..UsageRuntimeConfig::default()
});
- let plan = ExecutionPlan {
- request_id: "req-client-drop-cancels-upstream".into(),
- candidate_id: Some("cand-client-drop-cancels-upstream".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/chat/completions".into(),
- headers: BTreeMap::from([
- ("content-type".into(), "application/json".into()),
- ("accept".into(), "text/event-stream".into()),
- ]),
- content_type: Some("application/json".into()),
- content_encoding: None,
- body: RequestBody::from_json(json!({
- "model": "gpt-5.4",
- "messages": [],
- "stream": true
- })),
- stream: true,
- client_api_format: "openai:chat".into(),
- provider_api_format: "openai:chat".into(),
- model_name: Some("gpt-5.4".into()),
- proxy: None,
- transport_profile: None,
- timeouts: None,
- };
- let release_terminal = Arc::new(Notify::new());
- let terminal_frame_drained = Arc::new(Notify::new());
- let release_terminal_for_stream = Arc::clone(&release_terminal);
- let terminal_frame_drained_for_stream = Arc::clone(&terminal_frame_drained);
- let frame_stream = stream! {
+ let plan = ExecutionPlan {
+ request_id: "req-client-drop-cancels-upstream".into(),
+ candidate_id: Some("cand-client-drop-cancels-upstream".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/chat/completions".into(),
+ headers: BTreeMap::from([
+ ("content-type".into(), "application/json".into()),
+ ("accept".into(), "text/event-stream".into()),
+ ]),
+ content_type: Some("application/json".into()),
+ content_encoding: None,
+ body: RequestBody::from_json(json!({
+ "model": "gpt-5.4",
+ "messages": [],
+ "stream": true
+ })),
+ stream: true,
+ client_api_format: "openai:chat".into(),
+ provider_api_format: "openai:chat".into(),
+ model_name: Some("gpt-5.4".into()),
+ proxy: None,
+ transport_profile: None,
+ timeouts: None,
+ };
+ let release_terminal = Arc::new(Notify::new());
+ let terminal_frame_drained = Arc::new(Notify::new());
+ let release_terminal_for_stream = Arc::clone(&release_terminal);
+ let terminal_frame_drained_for_stream = Arc::clone(&terminal_frame_drained);
+ let frame_stream = stream! {
yield Ok::(Bytes::from_static(
b"{\"type\":\"headers\",\"payload\":{\"kind\":\"headers\",\"status_code\":200,\"headers\":{\"content-type\":\"text/event-stream\"}}}\n",
));
@@ -13669,119 +13671,157 @@ mod tests {
}
.boxed();
- let response = execute_stream_from_frame_stream(
- &state,
- plan,
- "trace-client-drop-cancels-upstream",
- &test_decision(),
- "openai_chat_stream",
- None,
- Some(json!({
- "request_id": "req-client-drop-cancels-upstream",
- "candidate_id": "cand-client-drop-cancels-upstream",
- "candidate_index": 0,
- "retry_index": 0,
- "provider_api_format": "openai:chat",
- "client_api_format": "openai:chat"
- })),
- crate::clock::current_unix_ms(),
- Instant::now(),
- RequestStageTrace::from_env(),
- true,
- frame_stream,
- None,
- )
- .await
- .expect("execution should succeed")
- .expect("execution should return a client response");
+ let response = crate::request_lifecycle::run_request(async move {
+ crate::request_lifecycle::configure_client_disconnect(
+ aether_routing_core::RoutingExecutionPolicy {
+ cancel_on_client_disconnect,
+ ..Default::default()
+ },
+ );
+ execute_stream_from_frame_stream(
+ &state,
+ plan,
+ "trace-client-drop-cancels-upstream",
+ &test_decision(),
+ "openai_chat_stream",
+ None,
+ Some(json!({
+ "request_id": "req-client-drop-cancels-upstream",
+ "candidate_id": "cand-client-drop-cancels-upstream",
+ "candidate_index": 0,
+ "retry_index": 0,
+ "provider_api_format": "openai:chat",
+ "client_api_format": "openai:chat"
+ })),
+ crate::clock::current_unix_ms(),
+ Instant::now(),
+ RequestStageTrace::from_env(),
+ true,
+ frame_stream,
+ None,
+ )
+ .await
+ .map(|response| response.expect("execution should return a client response"))
+ })
+ .await
+ .expect("execution should succeed");
- let mut body_stream = response.into_body().into_data_stream();
- let first = tokio::time::timeout(Duration::from_secs(1), async {
- loop {
- let chunk = body_stream
- .next()
- .await
- .expect("body should yield first chunk")
- .expect("first chunk should be ok");
- if chunk.as_ref() != b": aether-keepalive\n\n" {
- break chunk;
+ let mut body_stream = response.into_body().into_data_stream();
+ let first = tokio::time::timeout(Duration::from_secs(1), async {
+ loop {
+ let chunk = body_stream
+ .next()
+ .await
+ .expect("body should yield first chunk")
+ .expect("first chunk should be ok");
+ if chunk.as_ref() != b": aether-keepalive\n\n" {
+ break chunk;
+ }
}
- }
- })
- .await
- .expect("first business chunk should arrive");
- assert_eq!(
+ })
+ .await
+ .expect("first business chunk should arrive");
+ assert_eq!(
first.as_ref(),
b"data: {\"id\":\"first\",\"choices\":[{\"index\":0,\"delta\":{\"content\":\"hello\"}}]}\n\n"
);
- tokio::time::sleep(Duration::from_millis(30)).await;
- drop(body_stream);
- let candidates = tokio::time::timeout(Duration::from_secs(1), async {
- loop {
- let candidates = request_candidate_repository
- .list_by_request_id("req-client-drop-cancels-upstream")
- .await
- .expect("request candidates should read");
- if candidates
- .first()
- .is_some_and(|candidate| candidate.status == RequestCandidateStatus::Cancelled)
- {
- break candidates;
- }
- tokio::time::sleep(Duration::from_millis(10)).await;
+ tokio::time::sleep(Duration::from_millis(30)).await;
+ drop(body_stream);
+ if !cancel_on_client_disconnect {
+ release_terminal.notify_one();
}
- })
- .await
- .expect("candidate should be marked cancelled");
- assert_eq!(candidates[0].status_code, Some(499));
- assert_eq!(
- 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");
+ let expected_candidate_status = if cancel_on_client_disconnect {
+ RequestCandidateStatus::Cancelled
+ } else {
+ RequestCandidateStatus::Success
+ };
+ let expected_usage_status = if cancel_on_client_disconnect {
+ "cancelled"
+ } else {
+ "completed"
+ };
+ let candidates = tokio::time::timeout(Duration::from_secs(1), async {
+ loop {
+ let candidates = request_candidate_repository
+ .list_by_request_id("req-client-drop-cancels-upstream")
+ .await
+ .expect("request candidates should read");
+ if candidates
+ .first()
+ .is_some_and(|candidate| candidate.status == expected_candidate_status)
+ {
+ break candidates;
+ }
+ tokio::time::sleep(Duration::from_millis(10)).await;
}
- tokio::time::sleep(Duration::from_millis(10)).await;
- }
- })
- .await
- .expect("usage should be marked cancelled");
- assert_eq!(stored_usage.billing_status, "void");
- assert_eq!(stored_usage.status_code, Some(499));
- assert_eq!(stored_usage.input_tokens, 0);
- assert_eq!(stored_usage.output_tokens, 0);
- assert_eq!(stored_usage.total_tokens, 0);
- 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"
- );
-
- release_terminal.notify_one();
- assert!(
- tokio::time::timeout(
- Duration::from_millis(100),
- terminal_frame_drained.notified()
- )
+ })
.await
- .is_err(),
- "upstream frame stream should stop when the client disconnects"
- );
+ .expect("candidate should be marked cancelled");
+ assert_eq!(
+ candidates[0].status_code,
+ Some(if cancel_on_client_disconnect {
+ 499
+ } else {
+ 200
+ })
+ );
+ assert_eq!(
+ candidates[0].error_type.as_deref(),
+ cancel_on_client_disconnect.then_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 == expected_usage_status)
+ {
+ break usage.expect("cancelled usage should exist");
+ }
+ tokio::time::sleep(Duration::from_millis(10)).await;
+ }
+ })
+ .await
+ .expect("usage should be marked cancelled");
+ if !cancel_on_client_disconnect {
+ assert_ne!(stored_usage.billing_status, "void");
+ assert_eq!(stored_usage.status_code, Some(200));
+ assert_eq!(stored_usage.input_tokens, 7);
+ assert_eq!(stored_usage.output_tokens, 11);
+ assert_eq!(stored_usage.total_tokens, 18);
+ continue;
+ }
+ assert_eq!(stored_usage.billing_status, "void");
+ assert_eq!(stored_usage.status_code, Some(499));
+ assert_eq!(stored_usage.input_tokens, 0);
+ assert_eq!(stored_usage.output_tokens, 0);
+ assert_eq!(stored_usage.total_tokens, 0);
+ 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"
+ );
+
+ release_terminal.notify_one();
+ assert!(
+ tokio::time::timeout(
+ Duration::from_millis(100),
+ terminal_frame_drained.notified()
+ )
+ .await
+ .is_err(),
+ "upstream frame stream should stop when the client disconnects"
+ );
+ }
}
#[tokio::test]
diff --git a/apps/aether-gateway/src/executor/orchestration.rs b/apps/aether-gateway/src/executor/orchestration.rs
index 7fd6ebca6..e79f96038 100644
--- a/apps/aether-gateway/src/executor/orchestration.rs
+++ b/apps/aether-gateway/src/executor/orchestration.rs
@@ -822,15 +822,20 @@ where
let started_at = Instant::now();
let (tx, rx) = mpsc::channel::>(1);
let request_diagnostics = current_request_diagnostics();
+ let cancel_on_disconnect = crate::request_lifecycle::cancel_on_client_disconnect();
tokio::spawn(async move {
scope_request_diagnostics_with(request_diagnostics, async move {
- let bytes = standard_text_sync_heartbeat_final_bytes(
+ let completion = standard_text_sync_heartbeat_final_bytes(
client_api_format.as_str(),
redaction_slot.as_ref(),
- execute(state, parts, trace_id, decision, plan_kind, started_at).await,
- )
- .await;
+ tokio::select! {
+ biased;
+ _ = tx.closed(), if cancel_on_disconnect => return,
+ result = execute(state, parts, trace_id, decision, plan_kind, started_at) => result,
+ },
+ );
+ let bytes = completion.await;
let _ = tx.send(Ok(Bytes::from(bytes))).await;
})
.await;
@@ -1097,23 +1102,26 @@ fn build_openai_image_sync_heartbeat_shell_response(
let started_at = Instant::now();
let (tx, rx) = mpsc::channel::>(1);
let request_diagnostics = current_request_diagnostics();
+ let cancel_on_disconnect = crate::request_lifecycle::cancel_on_client_disconnect();
tokio::spawn(async move {
scope_request_diagnostics_with(request_diagnostics, async move {
- let bytes = openai_image_sync_heartbeat_final_bytes(
- execute_openai_image_sync_heartbeat_attempts(
- state,
- request_path,
- trace_id,
- decision,
- plan_kind,
- attempts,
- transfer_tracker,
- started_at,
- )
- .await,
- )
- .await;
+ let execution = execute_openai_image_sync_heartbeat_attempts(
+ state,
+ request_path,
+ trace_id,
+ decision,
+ plan_kind,
+ attempts,
+ transfer_tracker,
+ started_at,
+ );
+ let outcome = tokio::select! {
+ biased;
+ _ = tx.closed(), if cancel_on_disconnect => return,
+ result = execution => result,
+ };
+ let bytes = openai_image_sync_heartbeat_final_bytes(outcome).await;
let _ = tx.send(Ok(Bytes::from(bytes))).await;
})
.await;
@@ -2331,6 +2339,45 @@ mod tests {
.expect("background completion should release admission");
}
+ #[tokio::test]
+ async fn standard_text_sync_heartbeat_cancels_when_routing_policy_enables_it() {
+ let (started_tx, started_rx) = tokio::sync::oneshot::channel();
+ let (mut release_tx, release_rx) = tokio::sync::oneshot::channel::<()>();
+ let response = crate::request_lifecycle::run_request(async move {
+ crate::request_lifecycle::configure_client_disconnect(
+ aether_routing_core::RoutingExecutionPolicy {
+ cancel_on_client_disconnect: true,
+ ..Default::default()
+ },
+ );
+ let (parts, _) = http::Request::builder()
+ .method("POST")
+ .uri("/v1/responses")
+ .body(())
+ .unwrap()
+ .into_parts();
+ build_standard_text_sync_heartbeat_shell_response(
+ AppState::new().unwrap(),
+ parts,
+ "trace-heartbeat-disconnect".to_string(),
+ test_standard_text_heartbeat_decision(),
+ TEST_STANDARD_TEXT_SYNC_PLAN_KIND.to_string(),
+ move |_, _, _, _, _, _| async move {
+ started_tx.send(()).unwrap();
+ release_rx.await.unwrap();
+ Ok(LocalExecutionRequestOutcome::NoPath)
+ },
+ )
+ })
+ .await
+ .unwrap();
+ started_rx.await.unwrap();
+ drop(response);
+ tokio::time::timeout(Duration::from_secs(1), release_tx.closed())
+ .await
+ .expect("heartbeat must drop upstream execution immediately");
+ }
+
#[tokio::test]
async fn standard_text_sync_heartbeat_propagates_request_diagnostics_to_terminal_usage() {
let (state, usage_repository) = heartbeat_usage_test_state(json!({
diff --git a/apps/aether-gateway/src/handlers/proxy/mod.rs b/apps/aether-gateway/src/handlers/proxy/mod.rs
index 5d8886ce4..0f8f77fd8 100644
--- a/apps/aether-gateway/src/handlers/proxy/mod.rs
+++ b/apps/aether-gateway/src/handlers/proxy/mod.rs
@@ -1035,7 +1035,7 @@ pub(crate) async fn proxy_request(
ConnectInfo(remote_addr): ConnectInfo,
request: Request,
) -> Result, GatewayError> {
- crate::request_diagnostics::scope_request_diagnostics(Box::pin(proxy_request_inner(
+ crate::request_lifecycle::run_request(Box::pin(proxy_request_inner(
state,
remote_addr,
request,
diff --git a/apps/aether-gateway/src/handlers/proxy/websocket/responses/connection.rs b/apps/aether-gateway/src/handlers/proxy/websocket/responses/connection.rs
index 5749b714c..23774c5f5 100644
--- a/apps/aether-gateway/src/handlers/proxy/websocket/responses/connection.rs
+++ b/apps/aether-gateway/src/handlers/proxy/websocket/responses/connection.rs
@@ -68,7 +68,12 @@ pub(super) async fn relay_bound_connection(
state: &AppState,
context: &WebSocketRequestContext,
) {
+ let mut client_connected = true;
loop {
+ if !client_connected && !bound.turn_state.response_in_flight() {
+ close_bound_upstream(bound).await;
+ break;
+ }
let active_turn_deadline = bound.turn_state.attempt().map(|turn| turn.deadline());
tokio::select! {
_ = wait_for_optional_deadline(active_turn_deadline.map(|deadline| deadline.deadline)) => {
@@ -100,8 +105,12 @@ pub(super) async fn relay_bound_connection(
).await;
break;
}
- client_message = client_socket.next() => {
+ client_message = client_socket.next(), if client_connected => {
let Some(client_message) = client_message else {
+ if retain_disconnected_turn(bound) {
+ client_connected = false;
+ continue;
+ }
finalize_active_turn(
bound,
state,
@@ -111,6 +120,10 @@ pub(super) async fn relay_bound_connection(
break;
};
let Ok(client_message) = client_message else {
+ if retain_disconnected_turn(bound) {
+ client_connected = false;
+ continue;
+ }
warn!(
event_name = "responses_websocket_client_receive_failed",
log_type = "ops",
@@ -127,6 +140,12 @@ pub(super) async fn relay_bound_connection(
close_bound_upstream(bound).await;
break;
};
+ if matches!(client_message, AxumWsMessage::Close(_))
+ && retain_disconnected_turn(bound)
+ {
+ client_connected = false;
+ continue;
+ }
match Box::pin(forward_client_message(
client_message,
bound,
@@ -559,6 +578,7 @@ pub(super) async fn relay_bound_connection(
let mut relay_send_error = None;
let mut relay_serialization_failed = false;
match relay_directive {
+ _ if !client_connected => {}
Some(ResponsesWebSocketRelayDirective::ForwardOriginal) => {
let client_frame = match parsed_upstream_frame.as_ref().map(|frame| {
bound
@@ -673,6 +693,10 @@ pub(super) async fn relay_bound_connection(
break;
}
if let Some(error) = relay_send_error {
+ if terminal_outcome.is_none() && retain_disconnected_turn(bound) {
+ client_connected = false;
+ continue;
+ }
warn!(
event_name = "responses_websocket_client_send_failed",
log_type = "ops",
@@ -737,6 +761,20 @@ pub(super) async fn relay_bound_connection(
}
}
+fn retain_disconnected_turn(bound: &mut BoundResponsesConnection) -> bool {
+ if bound
+ .turn_state
+ .attempt()
+ .is_none_or(|attempt| attempt.cancel_on_client_disconnect())
+ {
+ return false;
+ }
+ bound
+ .turn_state
+ .record_client_delivery_aborted(CLIENT_DELIVERY_FAILED_REASON);
+ true
+}
+
struct PendingContinuationRegistration {
user_id: String,
api_key_id: String,
diff --git a/apps/aether-gateway/src/handlers/proxy/websocket/responses/turn.rs b/apps/aether-gateway/src/handlers/proxy/websocket/responses/turn.rs
index 4fa9d7b07..34b73d12b 100644
--- a/apps/aether-gateway/src/handlers/proxy/websocket/responses/turn.rs
+++ b/apps/aether-gateway/src/handlers/proxy/websocket/responses/turn.rs
@@ -845,6 +845,13 @@ fn websocket_auth_rejection_error(rejection: GatewayLocalAuthRejection) -> Gatew
}
impl ResponsesProviderAttempt {
+ pub(super) fn cancel_on_client_disconnect(&self) -> bool {
+ crate::orchestration::routing_execution_policy_from_report_context(
+ self.lifecycle.report_context(),
+ )
+ .is_some_and(|policy| policy.cancel_on_client_disconnect)
+ }
+
/// Releases all per-turn capacity before terminal persistence starts.
/// Provider-pool runtime tokens normally use an awaited removal. The
/// bounded wait prevents a broken runtime backend from stalling the relay;
diff --git a/apps/aether-gateway/src/lib.rs b/apps/aether-gateway/src/lib.rs
index 630717169..57bd703a2 100644
--- a/apps/aether-gateway/src/lib.rs
+++ b/apps/aether-gateway/src/lib.rs
@@ -71,6 +71,7 @@ mod rate_limit;
mod request_candidate_queue;
mod request_candidate_runtime;
mod request_diagnostics;
+mod request_lifecycle;
mod roles;
mod router;
mod routing;
diff --git a/apps/aether-gateway/src/request_lifecycle.rs b/apps/aether-gateway/src/request_lifecycle.rs
new file mode 100644
index 000000000..107b1f7ea
--- /dev/null
+++ b/apps/aether-gateway/src/request_lifecycle.rs
@@ -0,0 +1,336 @@
+use std::future::Future;
+use std::pin::Pin;
+use std::sync::atomic::{AtomicBool, Ordering};
+use std::sync::Arc;
+use std::task::{Context, Poll};
+
+use aether_routing_core::RoutingExecutionPolicy;
+use axum::body::{Body, Bytes, HttpBody};
+use http::Response;
+use http_body::{Frame, SizeHint};
+use http_body_util::BodyExt;
+
+use crate::request_diagnostics::{scope_request_diagnostics_with, RequestDiagnostics};
+use crate::GatewayError;
+
+tokio::task_local! {
+ static CANCEL_ON_CLIENT_DISCONNECT: Arc;
+}
+
+pub(crate) fn configure_client_disconnect(policy: RoutingExecutionPolicy) {
+ let _ = CANCEL_ON_CLIENT_DISCONNECT.try_with(|cancel| {
+ cancel.store(policy.cancel_on_client_disconnect, Ordering::Release);
+ });
+}
+
+pub(crate) fn cancel_on_client_disconnect() -> bool {
+ CANCEL_ON_CLIENT_DISCONNECT
+ .try_with(|cancel| cancel.load(Ordering::Acquire))
+ .unwrap_or(false)
+}
+
+pub(crate) async fn run_request(future: F) -> Result, GatewayError>
+where
+ F: Future
-
@@ -970,6 +988,9 @@ const cfHeartbeat = computed(() => (
const cyberContinueFailover = computed(() => (
draft.value?.config_json.default_policy.cyber_continue_failover ?? false
))
+const cancelOnClientDisconnect = computed(() => (
+ draft.value?.config_json.default_policy.cancel_on_client_disconnect ?? false
+))
interface ModelRow {
name: string
displayName: string
@@ -1306,7 +1327,7 @@ function updateKeepPriorityOnConversion(value: boolean): void {
}
function updateExecutionPolicy(
- field: 'enable_cf_heartbeat' | 'cyber_continue_failover',
+ field: 'enable_cf_heartbeat' | 'cyber_continue_failover' | 'cancel_on_client_disconnect',
value: boolean,
): void {
if (!draft.value) return