diff --git a/scripts/check-pingora-gateway-smoke.mjs b/scripts/check-pingora-gateway-smoke.mjs index a9b69e170..4428baf8c 100644 --- a/scripts/check-pingora-gateway-smoke.mjs +++ b/scripts/check-pingora-gateway-smoke.mjs @@ -383,6 +383,8 @@ async function startApiMock() { const state = { requests: [], releaseHold: undefined, + releaseAbortHold: undefined, + closedBeforeResponse: [], }; const server = http.createServer(async (request, response) => { let body; @@ -406,11 +408,20 @@ async function startApiMock() { headers: request.headers, body, }); + response.on('close', () => { + if (!response.writableEnded) { + state.closedBeforeResponse.push(request.url || ''); + } + }); if (request.url?.endsWith('/hold')) { await new Promise((resolve) => { state.releaseHold = resolve; }); + } else if (request.url?.endsWith('/abort-hold')) { + await new Promise((resolve) => { + state.releaseAbortHold = resolve; + }); } else if (request.url?.endsWith('/upstream-close')) { request.socket.destroy(); return; @@ -1006,6 +1017,7 @@ async function runSmokeCases( await expectChunkedLimit(baseUrl, '/api/upload'); await expectConcurrencyLimit(baseUrl, api); + await expectDownstreamAbortCancelsUpstream(baseUrl, api); const rateLimitHeaders = { 'X-Forwarded-For': '203.0.113.13' }; const beforeRateLimitRequestCount = api.state.requests.length; @@ -1695,6 +1707,27 @@ async function expectConcurrencyLimit(baseUrl, api) { } } +async function expectDownstreamAbortCancelsUpstream(baseUrl, api) { + console.log('[pingora-gateway-smoke] 下游中断传播到 API 上游'); + const route = '/api/abort-hold'; + const hold = await openRawHttpRequest(`${baseUrl}${route}`, { + 'X-Forwarded-For': '203.0.113.16', + }); + + try { + await waitForCondition( + () => typeof api.state.releaseAbortHold === 'function', + ); + hold.socket.destroy(); + await hold.done.catch(() => undefined); + await waitForCondition(() => + api.state.closedBeforeResponse.includes(route), + ); + } finally { + api.state.releaseAbortHold?.(); + } +} + async function expectWebSocketUpgrade(baseUrl, route, spacetime) { console.log('[pingora-gateway-smoke] SpacetimeDB WebSocket Upgrade'); const url = new URL(route, baseUrl); diff --git a/server-rs/crates/api-server/src/editor_agent/api.rs b/server-rs/crates/api-server/src/editor_agent/api.rs index e26bd6b69..65a043142 100644 --- a/server-rs/crates/api-server/src/editor_agent/api.rs +++ b/server-rs/crates/api-server/src/editor_agent/api.rs @@ -443,6 +443,30 @@ fn editor_agent_attachment_requests_match( mod tests { use super::*; use shared_contracts::editor_agent::{EditorAgentAttachmentRef, EditorAgentAttachmentSource}; + use std::sync::Arc; + use tokio::sync::{Mutex as AsyncMutex, Semaphore}; + + #[derive(Clone)] + struct HttpAbortTestState { + handler_lock: Arc>, + handler_started: Arc, + handler_dropped: Arc, + } + + struct HttpAbortDropNotice(Arc); + + impl Drop for HttpAbortDropNotice { + fn drop(&mut self) { + self.0.add_permits(1); + } + } + + async fn pending_http_abort_test_handler(State(state): State) { + let _handler_lock_guard = state.handler_lock.lock().await; + let _drop_notice = HttpAbortDropNotice(state.handler_dropped.clone()); + state.handler_started.add_permits(1); + std::future::pending::<()>().await; + } fn attachment(reference_id: impl Into) -> EditorAgentAttachmentRef { EditorAgentAttachmentRef { @@ -590,6 +614,69 @@ mod tests { "美术 Agent 规划失败:规划总时长已达到 18 分钟安全上限" ); } + + #[tokio::test] + async fn dropping_an_http_request_drops_the_handler_and_releases_its_lock() { + let state = HttpAbortTestState { + handler_lock: Arc::new(AsyncMutex::new(())), + handler_started: Arc::new(Semaphore::new(0)), + handler_dropped: Arc::new(Semaphore::new(0)), + }; + let router = axum::Router::new() + .route( + "/pending", + axum::routing::post(pending_http_abort_test_handler), + ) + .with_state(state.clone()); + let listener = tokio::net::TcpListener::bind("127.0.0.1:0") + .await + .expect("test listener should bind"); + let address = listener + .local_addr() + .expect("test listener should have an address"); + let server = tokio::spawn(async move { + axum::serve(listener, router) + .await + .expect("test server should run"); + }); + let request = tokio::spawn(async move { + reqwest::Client::builder() + .pool_max_idle_per_host(0) + .build() + .expect("test client should build") + .post(format!("http://{address}/pending")) + .send() + .await + }); + + tokio::time::timeout( + Duration::from_secs(2), + state.handler_started.clone().acquire_owned(), + ) + .await + .expect("handler should start before cancellation") + .expect("handler start semaphore should stay open") + .forget(); + + request.abort(); + let _ = request.await; + + tokio::time::timeout( + Duration::from_secs(2), + state.handler_dropped.clone().acquire_owned(), + ) + .await + .expect("HTTP cancellation should drop the handler") + .expect("handler drop semaphore should stay open") + .forget(); + let _released_lock = + tokio::time::timeout(Duration::from_secs(2), state.handler_lock.lock()) + .await + .expect("HTTP cancellation should release the handler lock"); + + server.abort(); + let _ = server.await; + } } fn editor_agent_system_prompt() -> &'static str { r#"