Skip to content

Commit a595fbc

Browse files
committed
fix: prevent owned VM fetch wakeup starvation
1 parent a1b00be commit a595fbc

3 files changed

Lines changed: 47 additions & 179 deletions

File tree

crates/native-sidecar/src/execution/coordinator.rs

Lines changed: 1 addition & 2 deletions
Original file line numberDiff line numberDiff line change
@@ -690,12 +690,11 @@ where
690690
payload: VmFetchRequest,
691691
) -> OwnedVmRouteFuture {
692692
let input = self.prepare_owned_vm_route(request);
693-
let bridge = self.bridge.clone();
694693
let max_frame_bytes = self.config.max_frame_bytes;
695694
Box::pin(async move {
696695
let input = input?;
697696
let response_json =
698-
dispatch_owned_vm_fetch(bridge, &input.vm_id, input.vm.clone(), payload).await?;
697+
dispatch_owned_vm_fetch(&input.vm_id, input.vm.clone(), payload).await?;
699698
let response = agentos_native_sidecar_core::respond(
700699
&input.request,
701700
ResponsePayload::VmFetchResult(VmFetchResponse { response_json }),

crates/native-sidecar/src/execution/javascript/http.rs

Lines changed: 26 additions & 164 deletions
Original file line numberDiff line numberDiff line change
@@ -1721,104 +1721,6 @@ where
17211721
settle_owned_fetch_sync_rpc::<B>(vm, process_id, child_path, request, response).await
17221722
}
17231723

1724-
/// Poll and service one event for the private VM-fetch loop. Protocol-router
1725-
/// process events use `service_owned_javascript_sync_rpc_request` after claiming
1726-
/// the concrete event, preserving the broker's one-consumer rule.
1727-
pub(crate) async fn service_owned_root_javascript_event<B>(
1728-
bridge: &SharedBridge<B>,
1729-
vm_id: &str,
1730-
vm: &crate::state::VmHandle,
1731-
process_id: &str,
1732-
preserve_http_wait: bool,
1733-
) -> Result<bool, SidecarError>
1734-
where
1735-
B: NativeSidecarBridge + Send + 'static,
1736-
BridgeError<B>: fmt::Debug + Send + Sync + 'static,
1737-
{
1738-
enum Turn {
1739-
Idle,
1740-
Serviced,
1741-
Response {
1742-
request: JavascriptSyncRpcRequest,
1743-
response: Result<JavascriptSyncRpcServiceResponse, SidecarError>,
1744-
},
1745-
}
1746-
1747-
let turn = vm.try_command("service owned VM fetch target", |vm| {
1748-
let socket_paths = build_javascript_socket_path_context(vm)?;
1749-
let dns = vm.dns.clone();
1750-
let kernel_readiness = Arc::clone(&vm.kernel_socket_readiness);
1751-
let capabilities = vm.capabilities.clone();
1752-
let VmState {
1753-
kernel,
1754-
active_processes,
1755-
..
1756-
} = vm;
1757-
let process = active_processes.get_mut(process_id).ok_or_else(|| {
1758-
SidecarError::InvalidState(format!("vm.fetch target process disappeared: {process_id}"))
1759-
})?;
1760-
let event = process
1761-
.execution
1762-
.try_poll_event()
1763-
.map_err(|error| SidecarError::Execution(error.to_string()))?;
1764-
let Some(event) = event else {
1765-
return Ok(Turn::Idle);
1766-
};
1767-
match event {
1768-
ActiveExecutionEvent::JavascriptSyncRpcRequest(request)
1769-
if preserve_http_wait && request.method == "net.http_wait" =>
1770-
{
1771-
process.queue_pending_execution_event(
1772-
ActiveExecutionEvent::JavascriptSyncRpcRequest(request),
1773-
)?;
1774-
Ok(Turn::Serviced)
1775-
}
1776-
ActiveExecutionEvent::JavascriptSyncRpcRequest(request) => {
1777-
let mut future = Box::pin(service_javascript_sync_rpc(
1778-
JavascriptSyncRpcServiceRequest {
1779-
bridge,
1780-
vm_id,
1781-
dns: &dns,
1782-
socket_paths: &socket_paths,
1783-
kernel,
1784-
kernel_readiness,
1785-
process,
1786-
sync_request: &request,
1787-
capabilities,
1788-
},
1789-
));
1790-
let mut context = Context::from_waker(Waker::noop());
1791-
let poll = future.as_mut().poll(&mut context);
1792-
drop(future);
1793-
match poll {
1794-
Poll::Ready(response) => Ok(Turn::Response { request, response }),
1795-
Poll::Pending => {
1796-
process.queue_pending_execution_event(
1797-
ActiveExecutionEvent::JavascriptSyncRpcRequest(request),
1798-
)?;
1799-
Ok(Turn::Serviced)
1800-
}
1801-
}
1802-
}
1803-
ActiveExecutionEvent::Exited(code) => Err(SidecarError::Execution(format!(
1804-
"vm.fetch target exited before responding (exit code {code})"
1805-
))),
1806-
other => {
1807-
process.queue_pending_execution_event(other)?;
1808-
Ok(Turn::Serviced)
1809-
}
1810-
}
1811-
})?;
1812-
match turn {
1813-
Turn::Idle => Ok(false),
1814-
Turn::Serviced => Ok(true),
1815-
Turn::Response { request, response } => {
1816-
settle_owned_fetch_sync_rpc::<B>(vm, process_id, &[], request, response).await?;
1817-
Ok(true)
1818-
}
1819-
}
1820-
}
1821-
18221724
fn owned_fetch_process_notify(
18231725
vm: &crate::state::VmHandle,
18241726
process_id: &str,
@@ -1836,7 +1738,7 @@ fn owned_fetch_process_notify(
18361738
}
18371739

18381740
async fn wait_for_owned_fetch_progress(
1839-
socket: &OwnedKernelFetchSocket,
1741+
readiness_notify: &Arc<tokio::sync::Notify>,
18401742
process_notify: &Arc<tokio::sync::Notify>,
18411743
deadline: Instant,
18421744
) -> Result<(), SidecarError> {
@@ -1848,8 +1750,14 @@ async fn wait_for_owned_fetch_progress(
18481750
}
18491751
tokio::time::timeout(remaining, async {
18501752
tokio::select! {
1851-
_ = socket.readiness_notify.notified() => {},
1852-
_ = process_notify.notified() => {},
1753+
_ = readiness_notify.notified() => {},
1754+
_ = tokio::time::sleep(Duration::from_millis(1)) => {
1755+
// The ordinary dispatcher exclusively owns guest execution
1756+
// events. Wake it so runnable microtask or WASM work advances;
1757+
// this fetch path only owns the kernel response socket.
1758+
process_notify.notify_one();
1759+
tokio::task::yield_now().await;
1760+
},
18531761
}
18541762
})
18551763
.await
@@ -1862,9 +1770,7 @@ async fn wait_for_owned_fetch_progress(
18621770
Ok(())
18631771
}
18641772

1865-
async fn dispatch_owned_kernel_http_fetch<B>(
1866-
bridge: &SharedBridge<B>,
1867-
vm_id: &str,
1773+
async fn dispatch_owned_kernel_http_fetch(
18681774
vm: &crate::state::VmHandle,
18691775
target_process_id: &str,
18701776
port: u16,
@@ -1873,11 +1779,7 @@ async fn dispatch_owned_kernel_http_fetch<B>(
18731779
headers: &HttpHeaderCollection,
18741780
body_bytes: Option<&[u8]>,
18751781
max_fetch_response_bytes: usize,
1876-
) -> Result<String, SidecarError>
1877-
where
1878-
B: NativeSidecarBridge + Send + 'static,
1879-
BridgeError<B>: fmt::Debug + Send + Sync + 'static,
1880-
{
1782+
) -> Result<String, SidecarError> {
18811783
let mut socket = open_owned_kernel_fetch_socket(
18821784
vm,
18831785
target_process_id,
@@ -1911,23 +1813,20 @@ where
19111813
preview.chars().take(200).collect::<String>()
19121814
)));
19131815
}
1914-
let serviced =
1915-
service_owned_root_javascript_event(bridge, vm_id, vm, target_process_id, true).await?;
19161816
let progressed = poll_owned_kernel_fetch_socket(
19171817
&socket,
19181818
&mut response_buffer,
19191819
&mut peer_closed,
19201820
"vm.fetch",
19211821
)?;
1922-
if !progressed && !serviced {
1923-
wait_for_owned_fetch_progress(&socket, &process_notify, deadline).await?;
1822+
if !progressed {
1823+
wait_for_owned_fetch_progress(&socket.readiness_notify, &process_notify, deadline)
1824+
.await?;
19241825
}
19251826
}
19261827
}
19271828

1928-
async fn start_owned_kernel_http_fetch_stream<B>(
1929-
bridge: &SharedBridge<B>,
1930-
vm_id: &str,
1829+
async fn start_owned_kernel_http_fetch_stream(
19311830
vm: &crate::state::VmHandle,
19321831
target_process_id: &str,
19331832
port: u16,
@@ -1936,11 +1835,7 @@ async fn start_owned_kernel_http_fetch_stream<B>(
19361835
headers: &HttpHeaderCollection,
19371836
body_bytes: Option<&[u8]>,
19381837
max_response_bytes: usize,
1939-
) -> Result<String, SidecarError>
1940-
where
1941-
B: NativeSidecarBridge + Send + 'static,
1942-
BridgeError<B>: fmt::Debug + Send + Sync + 'static,
1943-
{
1838+
) -> Result<String, SidecarError> {
19441839
let socket = open_owned_kernel_fetch_socket(
19451840
vm,
19461841
target_process_id,
@@ -1981,16 +1876,15 @@ where
19811876
http_loopback_request_timeout().as_millis()
19821877
)));
19831878
}
1984-
let serviced =
1985-
service_owned_root_javascript_event(bridge, vm_id, vm, target_process_id, true).await?;
19861879
let progressed = poll_owned_kernel_fetch_socket(
19871880
&socket,
19881881
&mut response_buffer,
19891882
&mut peer_closed,
19901883
"vm.fetchStream",
19911884
)?;
1992-
if !progressed && !serviced {
1993-
wait_for_owned_fetch_progress(&socket, &process_notify, deadline).await?;
1885+
if !progressed {
1886+
wait_for_owned_fetch_progress(&socket.readiness_notify, &process_notify, deadline)
1887+
.await?;
19941888
}
19951889
};
19961890

@@ -2137,17 +2031,11 @@ fn lease_owned_fetch_stream(
21372031
})
21382032
}
21392033

2140-
async fn read_owned_kernel_http_fetch_stream<B>(
2141-
bridge: &SharedBridge<B>,
2142-
vm_id: &str,
2034+
async fn read_owned_kernel_http_fetch_stream(
21432035
vm: &crate::state::VmHandle,
21442036
stream_id: &str,
21452037
requested_max_bytes: usize,
2146-
) -> Result<String, SidecarError>
2147-
where
2148-
B: NativeSidecarBridge + Send + 'static,
2149-
BridgeError<B>: fmt::Debug + Send + Sync + 'static,
2150-
{
2038+
) -> Result<String, SidecarError> {
21512039
let max_bytes = requested_max_bytes.clamp(1, VM_FETCH_STREAM_CHUNK_MAX_BYTES);
21522040
let mut lease = lease_owned_fetch_stream(vm, stream_id)?;
21532041
let target_process_id = lease.state().target_process_id.clone();
@@ -2166,9 +2054,6 @@ where
21662054
http_loopback_request_timeout().as_millis()
21672055
)));
21682056
}
2169-
let serviced =
2170-
service_owned_root_javascript_event(bridge, vm_id, vm, &target_process_id, true)
2171-
.await?;
21722057
let (kernel_pid, socket_id) = (lease.state().kernel_pid, lease.state().socket_id);
21732058
let mut raw_buffer = std::mem::take(&mut lease.state_mut().raw_buffer);
21742059
let mut peer_closed = lease.state().peer_closed;
@@ -2231,21 +2116,9 @@ where
22312116
lease.state_mut().peer_closed = peer_closed;
22322117
if progressed {
22332118
lease.state_mut().last_progress_at = Instant::now();
2234-
} else if !serviced {
2235-
let remaining = deadline.saturating_duration_since(Instant::now());
2236-
tokio::time::timeout(remaining, async {
2237-
tokio::select! {
2238-
_ = lease.readiness_notify.notified() => {},
2239-
_ = process_notify.notified() => {},
2240-
}
2241-
})
2242-
.await
2243-
.map_err(|_| {
2244-
SidecarError::Execution(format!(
2245-
"ERR_AGENTOS_VM_FETCH_TIMEOUT: stream produced no data for {} ms; raise AGENTOS_HTTP_LOOPBACK_REQUEST_TIMEOUT_MS",
2246-
http_loopback_request_timeout().as_millis()
2247-
))
2248-
})?;
2119+
} else {
2120+
wait_for_owned_fetch_progress(&lease.readiness_notify, &process_notify, deadline)
2121+
.await?;
22492122
}
22502123
}
22512124

@@ -2358,16 +2231,11 @@ async fn dispatch_owned_loopback_http_request(
23582231

23592232
/// Run `vm.fetch` without holding either the process coordinator or a VM state
23602233
/// borrow over readiness and adapter-response waits.
2361-
pub(in crate::execution) async fn dispatch_owned_vm_fetch<B>(
2362-
bridge: SharedBridge<B>,
2234+
pub(in crate::execution) async fn dispatch_owned_vm_fetch(
23632235
vm_id: &str,
23642236
vm: crate::state::VmHandle,
23652237
payload: VmFetchRequest,
2366-
) -> Result<String, SidecarError>
2367-
where
2368-
B: NativeSidecarBridge + Send + 'static,
2369-
BridgeError<B>: fmt::Debug + Send + Sync + 'static,
2370-
{
2238+
) -> Result<String, SidecarError> {
23712239
let stream_operation = payload.stream_operation.as_deref();
23722240
if matches!(stream_operation, Some("read" | "cancel")) {
23732241
let stream_id = payload.stream_id.as_deref().ok_or_else(|| {
@@ -2377,8 +2245,6 @@ where
23772245
})?;
23782246
return if stream_operation == Some("read") {
23792247
read_owned_kernel_http_fetch_stream(
2380-
&bridge,
2381-
vm_id,
23822248
&vm,
23832249
stream_id,
23842250
payload.max_bytes.unwrap_or(64 * 1024) as usize,
@@ -2440,8 +2306,6 @@ where
24402306
if let Some(target_process_id) = kernel_target {
24412307
return if stream_operation == Some("start") {
24422308
start_owned_kernel_http_fetch_stream(
2443-
&bridge,
2444-
vm_id,
24452309
&vm,
24462310
&target_process_id,
24472311
payload.port,
@@ -2454,8 +2318,6 @@ where
24542318
.await
24552319
} else {
24562320
dispatch_owned_kernel_http_fetch(
2457-
&bridge,
2458-
vm_id,
24592321
&vm,
24602322
&target_process_id,
24612323
payload.port,

packages/core/tests/network-http-request.test.ts

Lines changed: 20 additions & 13 deletions
Original file line numberDiff line numberDiff line change
@@ -320,20 +320,27 @@ describe("guest http.request transport", () => {
320320
),
321321
]);
322322

323-
const head = await Promise.race([
324-
vm.fetchStreamStart(3000, new Request("http://guest/events")),
325-
new Promise<never>((_, reject) =>
326-
setTimeout(
327-
() => reject(new Error("stream response head timed out")),
328-
5_000,
323+
for (let requestIndex = 0; requestIndex < 10; requestIndex++) {
324+
const head = await Promise.race([
325+
vm.fetchStreamStart(
326+
3000,
327+
new Request(`http://guest/events-${requestIndex}`),
329328
),
330-
),
331-
]);
332-
const chunk = await vm.fetchStreamRead(head.streamId);
333-
334-
expect(head.status, stderr).toBe(200);
335-
expect(textDecoder.decode(chunk.body)).toBe("data: websocket-3\n\n");
336-
await vm.fetchStreamCancel(head.streamId);
329+
new Promise<never>((_, reject) =>
330+
setTimeout(
331+
() => reject(new Error("stream response head timed out")),
332+
5_000,
333+
),
334+
),
335+
]);
336+
const chunk = await vm.fetchStreamRead(head.streamId);
337+
338+
expect(head.status, stderr).toBe(200);
339+
expect(textDecoder.decode(chunk.body)).toBe(
340+
"data: websocket-3\n\n",
341+
);
342+
await vm.fetchStreamCancel(head.streamId);
343+
}
337344
await vm.process.kill(child.pid, "SIGKILL");
338345
} finally {
339346
upstreamWebSocket.clients.forEach((socket) => socket.terminate());

0 commit comments

Comments
 (0)