Merge origin/main into fix/gemini-cli-v1internal

This commit is contained in:
Mas0nShi
2026-05-28 11:58:00 +08:00
410 changed files with 38026 additions and 6621 deletions
File diff suppressed because it is too large Load Diff
@@ -1,4 +1,5 @@
use std::collections::BTreeMap;
use std::future::Future;
use std::io::Error as IoError;
use std::net::IpAddr;
use std::sync::OnceLock;
@@ -29,7 +30,9 @@ use crate::execution_runtime::ndjson::encode_stream_frame_ndjson;
use crate::execution_runtime::transport::{
build_browser_wreq_client, build_request_body, build_request_headers,
decode_response_body_bytes, format_upstream_request_error, format_wreq_upstream_request_error,
send_request, DirectHttpResponse, ExecutionRuntimeTransportError, ExecutionTransportControls,
resolve_stream_first_byte_timeout, send_request, stream_first_byte_timeout_message,
with_non_stream_total_timeout, DirectHttpResponse, ExecutionRuntimeTransportError,
ExecutionTransportControls,
};
const GROK_INTERNAL_HEADER: &str = "x-aether-grok-runtime";
@@ -148,9 +151,12 @@ pub(crate) async fn maybe_execute_grok_sync(
if !is_grok_plan(plan, report_context) {
return Ok(None);
}
let mut collected = execute_grok_app_chat(plan, report_context).await?;
materialize_grok_image_assets(plan, &mut collected).await;
Ok(Some(grok_execution_result(plan, collected, report_context)))
with_non_stream_total_timeout(plan, async move {
let mut collected = execute_grok_app_chat(plan, report_context).await?;
materialize_grok_image_assets(plan, &mut collected).await;
Ok(Some(grok_execution_result(plan, collected, report_context)))
})
.await
}
pub(crate) async fn maybe_execute_grok_stream(
@@ -387,6 +393,7 @@ async fn grok_imagine_websocket_images(
plan.proxy.as_ref(),
profile,
ExecutionTransportControls::default(),
true,
)?;
let response = client
.websocket(GROK_IMAGINE_WS_URL)
@@ -561,6 +568,7 @@ fn grok_success_frame_stream(
started_at: Instant,
mut body_stream: GrokUpstreamBodyStream,
) -> BoxStream<'static, Result<Bytes, IoError>> {
let stream_first_byte_timeout = resolve_stream_first_byte_timeout(&plan);
async_stream::stream! {
match encode_grok_headers_frame(
status_code,
@@ -581,8 +589,36 @@ fn grok_success_frame_stream(
let mut text_len = 0usize;
let mut thinking_len = 0usize;
let mut image_len = 0usize;
let mut terminal_error_emitted = false;
while let Some(item) = body_stream.next().await {
loop {
let item = if ttfb_ms.is_none() {
match await_grok_stream_first_byte(
body_stream.next(),
started_at,
stream_first_byte_timeout,
)
.await
{
Ok(item) => item,
Err(timeout) => {
match encode_grok_first_byte_timeout_frame(timeout) {
Ok(frame) => yield Ok(frame),
Err(err) => {
yield Err(err);
return;
}
}
terminal_error_emitted = true;
break;
}
}
} else {
body_stream.next().await
};
let Some(item) = item else {
break;
};
let chunk = match item {
Ok(chunk) => chunk,
Err(message) => {
@@ -593,6 +629,7 @@ fn grok_success_frame_stream(
return;
}
}
terminal_error_emitted = true;
break;
}
};
@@ -630,6 +667,25 @@ fn grok_success_frame_stream(
}
}
if terminal_error_emitted {
match encode_grok_telemetry_frame(
ttfb_ms,
Some(started_at.elapsed().as_millis() as u64),
upstream_bytes,
) {
Ok(frame) => yield Ok(frame),
Err(err) => {
yield Err(err);
return;
}
}
match encode_stream_frame_ndjson(&StreamFrame::eof_with_summary(None)) {
Ok(frame) => yield Ok(frame),
Err(err) => yield Err(err),
}
return;
}
adapter.finish();
match emit_grok_adapter_deltas(
&mut client_emitter,
@@ -691,6 +747,28 @@ fn grok_success_frame_stream(
.boxed()
}
async fn await_grok_stream_first_byte<T, F>(
future: F,
started_at: Instant,
timeout: Option<Duration>,
) -> Result<T, Duration>
where
F: Future<Output = T>,
{
let Some(timeout) = timeout else {
return Ok(future.await);
};
let Some(remaining) = timeout.checked_sub(started_at.elapsed()) else {
return Err(timeout);
};
if remaining.is_zero() {
return Err(timeout);
}
tokio::time::timeout(remaining, future)
.await
.map_err(|_| timeout)
}
fn emit_grok_adapter_deltas(
client_emitter: &mut GrokClientStreamEmitter,
adapter: &GrokStreamAdapter,
@@ -790,6 +868,22 @@ fn encode_grok_error_frame(status_code: u16, message: String) -> Result<Bytes, I
})
}
fn encode_grok_first_byte_timeout_frame(timeout: Duration) -> Result<Bytes, IoError> {
encode_stream_frame_ndjson(&StreamFrame {
frame_type: StreamFrameType::Error,
payload: StreamFramePayload::Error {
error: aether_contracts::ExecutionError {
kind: aether_contracts::ExecutionErrorKind::FirstByteTimeout,
phase: aether_contracts::ExecutionPhase::FirstByte,
message: stream_first_byte_timeout_message(timeout),
upstream_status: Some(504),
retryable: true,
failover_recommended: true,
},
},
})
}
enum GrokClientStreamEmitter {
OpenAiChat {
id: String,
File diff suppressed because it is too large Load Diff
@@ -1,6 +1,7 @@
use std::collections::BTreeMap;
use std::future::Future;
use std::io::Error as IoError;
use std::time::Instant;
use std::time::{Duration, Instant};
use aether_contracts::{
ExecutionError, ExecutionErrorKind, ExecutionPhase, ExecutionStreamTerminalSummary,
@@ -19,7 +20,7 @@ use crate::ai_serving::api::{
};
use crate::execution_runtime::ndjson::encode_stream_frame_ndjson;
use crate::execution_runtime::transport::{
format_wreq_upstream_request_error, DirectUpstreamResponse,
format_wreq_upstream_request_error, stream_first_byte_timeout_message, DirectUpstreamResponse,
};
use crate::execution_runtime::DirectUpstreamStreamExecution;
use crate::GatewayError;
@@ -37,6 +38,7 @@ pub(crate) fn build_direct_execution_frame_stream(
stream_summary_report_context,
response,
started_at,
stream_first_byte_timeout,
} = execution;
let mut observer_context = stream_summary_report_context;
@@ -64,7 +66,7 @@ pub(crate) fn build_direct_execution_frame_stream(
if should_buffer_non_stream_response(&headers, &observer_context) {
let original_headers = headers.clone();
match buffer_non_sse_upstream_body(response, started_at).await {
match buffer_non_sse_upstream_body(response, started_at, stream_first_byte_timeout).await {
Ok(buffered) => {
let mut response_headers = original_headers;
let mut response_body = Bytes::from(buffered.body_bytes);
@@ -134,6 +136,7 @@ pub(crate) fn build_direct_execution_frame_stream(
message,
ttfb_ms,
upstream_bytes,
first_byte_timeout,
}) => {
match encode_headers_frame(status_code, original_headers) {
Ok(frame) => yield Ok(frame),
@@ -142,7 +145,12 @@ pub(crate) fn build_direct_execution_frame_stream(
return;
}
}
match encode_error_frame(status_code, message) {
let error_frame = if let Some(timeout) = first_byte_timeout {
encode_first_byte_timeout_frame(timeout)
} else {
encode_error_frame(status_code, message)
};
match error_frame {
Ok(frame) => yield Ok(frame),
Err(err) => {
yield Err(err);
@@ -183,7 +191,33 @@ pub(crate) fn build_direct_execution_frame_stream(
match response {
DirectUpstreamResponse::Reqwest(response) => {
let mut bytes_stream = response.bytes_stream();
while let Some(item) = bytes_stream.next().await {
loop {
let item = if ttfb_ms.is_none() {
match await_stream_first_byte(
bytes_stream.next(),
started_at,
stream_first_byte_timeout,
)
.await
{
Ok(item) => item,
Err(timeout) => {
match encode_first_byte_timeout_frame(timeout) {
Ok(frame) => yield Ok(frame),
Err(err) => {
yield Err(err);
return;
}
}
break;
}
}
} else {
bytes_stream.next().await
};
let Some(item) = item else {
break;
};
match item {
Ok(chunk) => {
if ttfb_ms.is_none() {
@@ -239,7 +273,33 @@ pub(crate) fn build_direct_execution_frame_stream(
}
DirectUpstreamResponse::BrowserWreq(response) => {
let mut bytes_stream = response.bytes_stream();
while let Some(item) = bytes_stream.next().await {
loop {
let item = if ttfb_ms.is_none() {
match await_stream_first_byte(
bytes_stream.next(),
started_at,
stream_first_byte_timeout,
)
.await
{
Ok(item) => item,
Err(timeout) => {
match encode_first_byte_timeout_frame(timeout) {
Ok(frame) => yield Ok(frame),
Err(err) => {
yield Err(err);
return;
}
}
break;
}
}
} else {
bytes_stream.next().await
};
let Some(item) = item else {
break;
};
match item {
Ok(chunk) => {
if ttfb_ms.is_none() {
@@ -294,7 +354,30 @@ pub(crate) fn build_direct_execution_frame_stream(
}
}
DirectUpstreamResponse::LocalTunnel(mut response) => loop {
match response.next_chunk().await {
let item = if ttfb_ms.is_none() {
match await_stream_first_byte(
response.next_chunk(),
started_at,
stream_first_byte_timeout,
)
.await
{
Ok(item) => item,
Err(timeout) => {
match encode_first_byte_timeout_frame(timeout) {
Ok(frame) => yield Ok(frame),
Err(err) => {
yield Err(err);
return;
}
}
break;
}
}
} else {
response.next_chunk().await
};
match item {
Ok(Some(chunk)) => {
if ttfb_ms.is_none() {
ttfb_ms = Some(started_at.elapsed().as_millis() as u64);
@@ -428,6 +511,44 @@ fn encode_error_frame(status_code: u16, message: String) -> Result<Bytes, IoErro
})
}
fn encode_first_byte_timeout_frame(timeout: Duration) -> Result<Bytes, IoError> {
encode_stream_frame_ndjson(&StreamFrame {
frame_type: StreamFrameType::Error,
payload: StreamFramePayload::Error {
error: ExecutionError {
kind: ExecutionErrorKind::FirstByteTimeout,
phase: ExecutionPhase::FirstByte,
message: stream_first_byte_timeout_message(timeout),
upstream_status: Some(504),
retryable: true,
failover_recommended: true,
},
},
})
}
async fn await_stream_first_byte<T, F>(
future: F,
started_at: Instant,
timeout: Option<Duration>,
) -> Result<T, Duration>
where
F: Future<Output = T>,
{
let Some(timeout) = timeout else {
return Ok(future.await);
};
let Some(remaining) = timeout.checked_sub(started_at.elapsed()) else {
return Err(timeout);
};
if remaining.is_zero() {
return Err(timeout);
}
tokio::time::timeout(remaining, future)
.await
.map_err(|_| timeout)
}
struct BufferedUpstreamBody {
body_bytes: Vec<u8>,
ttfb_ms: Option<u64>,
@@ -438,6 +559,7 @@ struct BufferedUpstreamBodyError {
message: String,
ttfb_ms: Option<u64>,
upstream_bytes: u64,
first_byte_timeout: Option<Duration>,
}
fn response_headers_indicate_sse(headers: &BTreeMap<String, String>) -> bool {
@@ -477,6 +599,7 @@ fn should_buffer_non_stream_response(
async fn buffer_non_sse_upstream_body(
response: DirectUpstreamResponse,
started_at: Instant,
stream_first_byte_timeout: Option<Duration>,
) -> Result<BufferedUpstreamBody, BufferedUpstreamBodyError> {
let mut body_bytes = Vec::new();
let mut upstream_bytes = 0u64;
@@ -485,7 +608,31 @@ async fn buffer_non_sse_upstream_body(
match response {
DirectUpstreamResponse::Reqwest(response) => {
let mut bytes_stream = response.bytes_stream();
while let Some(item) = bytes_stream.next().await {
loop {
let item = if ttfb_ms.is_none() {
match await_stream_first_byte(
bytes_stream.next(),
started_at,
stream_first_byte_timeout,
)
.await
{
Ok(item) => item,
Err(timeout) => {
return Err(BufferedUpstreamBodyError {
message: stream_first_byte_timeout_message(timeout),
ttfb_ms,
upstream_bytes,
first_byte_timeout: Some(timeout),
});
}
}
} else {
bytes_stream.next().await
};
let Some(item) = item else {
break;
};
match item {
Ok(chunk) => {
if ttfb_ms.is_none() {
@@ -507,6 +654,7 @@ async fn buffer_non_sse_upstream_body(
message,
ttfb_ms,
upstream_bytes,
first_byte_timeout: None,
});
}
}
@@ -514,7 +662,31 @@ async fn buffer_non_sse_upstream_body(
}
DirectUpstreamResponse::BrowserWreq(response) => {
let mut bytes_stream = response.bytes_stream();
while let Some(item) = bytes_stream.next().await {
loop {
let item = if ttfb_ms.is_none() {
match await_stream_first_byte(
bytes_stream.next(),
started_at,
stream_first_byte_timeout,
)
.await
{
Ok(item) => item,
Err(timeout) => {
return Err(BufferedUpstreamBodyError {
message: stream_first_byte_timeout_message(timeout),
ttfb_ms,
upstream_bytes,
first_byte_timeout: Some(timeout),
});
}
}
} else {
bytes_stream.next().await
};
let Some(item) = item else {
break;
};
match item {
Ok(chunk) => {
if ttfb_ms.is_none() {
@@ -536,13 +708,35 @@ async fn buffer_non_sse_upstream_body(
message,
ttfb_ms,
upstream_bytes,
first_byte_timeout: None,
});
}
}
}
}
DirectUpstreamResponse::LocalTunnel(mut response) => loop {
match response.next_chunk().await {
let item = if ttfb_ms.is_none() {
match await_stream_first_byte(
response.next_chunk(),
started_at,
stream_first_byte_timeout,
)
.await
{
Ok(item) => item,
Err(timeout) => {
return Err(BufferedUpstreamBodyError {
message: stream_first_byte_timeout_message(timeout),
ttfb_ms,
upstream_bytes,
first_byte_timeout: Some(timeout),
});
}
}
} else {
response.next_chunk().await
};
match item {
Ok(Some(chunk)) => {
if ttfb_ms.is_none() {
ttfb_ms = Some(started_at.elapsed().as_millis() as u64);
@@ -563,6 +757,7 @@ async fn buffer_non_sse_upstream_body(
message,
ttfb_ms,
upstream_bytes,
first_byte_timeout: None,
});
}
}
@@ -761,6 +956,7 @@ mod tests {
use base64::Engine as _;
use futures_util::StreamExt;
use serde_json::Value;
use tokio::io::{AsyncReadExt, AsyncWriteExt};
use tokio::sync::watch;
use super::{
@@ -903,6 +1099,93 @@ mod tests {
);
}
#[tokio::test]
async fn direct_execution_frame_stream_applies_first_byte_timeout_after_headers() {
let listener = crate::test_support::bind_loopback_listener()
.await
.expect("listener should bind");
let addr = listener.local_addr().expect("local addr should resolve");
let server = tokio::spawn(async move {
let (mut socket, _) = listener.accept().await.expect("client should connect");
let mut request = [0_u8; 1024];
let _ = socket
.read(&mut request)
.await
.expect("request should read");
socket
.write_all(
b"HTTP/1.1 200 OK\r\ncontent-type: text/event-stream\r\ntransfer-encoding: chunked\r\n\r\n",
)
.await
.expect("headers should write");
socket.flush().await.expect("headers should flush");
tokio::time::sleep(Duration::from_millis(200)).await;
let _ = socket.write_all(b"d\r\ndata: hello\n\n\r\n0\r\n\r\n").await;
});
let execution = DirectSyncExecutionRuntime::new()
.execute_stream(&ExecutionPlan {
request_id: "req-stream-first-byte-timeout".into(),
candidate_id: Some("cand-stream-first-byte-timeout".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: format!("http://{addr}/chat"),
headers: BTreeMap::from([("content-type".into(), "application/json".into())]),
content_type: Some("application/json".into()),
content_encoding: None,
body: RequestBody::from_json(serde_json::json!({"stream": true})),
stream: true,
client_api_format: "openai:chat".into(),
provider_api_format: "openai:chat".into(),
model_name: Some("gpt-5".into()),
proxy: None,
transport_profile: None,
timeouts: Some(ExecutionTimeouts {
first_byte_ms: Some(50),
total_ms: Some(5_000),
..ExecutionTimeouts::default()
}),
})
.await
.expect("stream execution should receive response headers");
let frames = build_direct_execution_frame_stream(execution)
.map(|item| item.expect("frame should encode"))
.collect::<Vec<_>>()
.await
.into_iter()
.map(|bytes| String::from_utf8(bytes.to_vec()).expect("frame should be utf8"))
.collect::<Vec<_>>();
server.abort();
let error_frame = frames
.iter()
.map(|line| serde_json::from_str::<Value>(line).expect("frame should parse"))
.find(|frame| frame.get("type").and_then(Value::as_str) == Some("error"))
.expect("timeout should emit an error frame");
assert_eq!(
error_frame
.get("payload")
.and_then(|payload| payload.get("error"))
.and_then(|error| error.get("kind"))
.and_then(Value::as_str),
Some("first_byte_timeout")
);
assert!(error_frame
.get("payload")
.and_then(|payload| payload.get("error"))
.and_then(|error| error.get("message"))
.and_then(Value::as_str)
.is_some_and(
|message| message.contains("provider stream first byte timeout after 50 ms")
));
}
#[tokio::test]
async fn direct_execution_frame_stream_emits_telemetry_before_first_data_frame() {
let listener = crate::test_support::bind_loopback_listener()
File diff suppressed because it is too large Load Diff
@@ -1,6 +1,6 @@
use std::collections::{BTreeMap, HashMap, HashSet};
use std::fs;
use std::io::Error as IoError;
use std::io::{Error as IoError, Read, Seek, SeekFrom};
use std::net::{SocketAddr, TcpListener, TcpStream};
use std::path::{Path, PathBuf};
use std::process::{Child, Command, Stdio};
@@ -18,9 +18,9 @@ use aether_provider_transport::windsurf::cascade::{
build_get_trajectory_steps_request, build_get_user_status_request, build_heartbeat_request,
build_initialize_panel_state_request, build_send_cascade_message_request_with_options,
build_start_cascade_request, build_update_panel_state_with_user_status_request,
build_update_workspace_trust_request, extract_grpc_frames, extract_user_status_bytes,
grpc_frame, parse_generator_metadata, parse_start_cascade_response, parse_trajectory_status,
parse_trajectory_steps, CascadeImage, CascadeUsage, SendCascadeMessageOptions,
build_update_workspace_trust_request, extract_user_status_bytes, parse_generator_metadata,
parse_start_cascade_response, parse_trajectory_status, parse_trajectory_steps, CascadeImage,
CascadeUsage, SendCascadeMessageOptions,
};
use aether_provider_transport::windsurf::models::resolve_windsurf_model;
use aether_provider_transport::windsurf::{GET_CHAT_MESSAGE_PATH, WINDSURF_ENVELOPE_NAME};
@@ -32,11 +32,11 @@ use regex::Regex;
use serde_json::{json, Value};
use sha2::{Digest, Sha256};
use tokio::sync::mpsc;
use tracing::{debug, error, info, warn};
use tracing::{debug, info, warn};
use uuid::Uuid;
use super::ndjson::encode_stream_frame_ndjson;
use super::transport::ExecutionRuntimeTransportError;
use super::transport::{with_non_stream_total_timeout, ExecutionRuntimeTransportError};
use crate::AppState;
const LS_SERVICE: &str = "/exa.language_server_pb.LanguageServerService";
@@ -51,6 +51,7 @@ const CASCADE_TEXT_STALL: Duration = Duration::from_secs(45);
const CASCADE_THINKING_STALL: Duration = Duration::from_secs(120);
const SSE_HEARTBEAT_INTERVAL: Duration = Duration::from_secs(15);
const LS_READY_TIMEOUT: Duration = Duration::from_secs(25);
const WINDOWS_LS_READY_TIMEOUT: Duration = Duration::from_secs(90);
const GRPC_SHORT_TIMEOUT: Duration = Duration::from_secs(5);
const GRPC_STATUS_TIMEOUT: Duration = Duration::from_secs(10);
const GRPC_REQUEST_TIMEOUT: Duration = Duration::from_secs(45);
@@ -215,66 +216,69 @@ pub(crate) async fn maybe_execute_windsurf_sync(
let Some(input) = detect_windsurf_request(plan, report_context) else {
return Ok(None);
};
let key_upstream_metadata = read_windsurf_key_upstream_metadata(state, plan).await;
let prepared = prepare_windsurf_cascade(plan, input, key_upstream_metadata).await?;
let started_at = Instant::now();
let mut deltas = Vec::new();
let poll_result = poll_windsurf_cascade_with_transport_recovery(&prepared, |event| {
if let WindsurfPollEvent::TextDelta(delta) = event {
deltas.push(sanitize_windsurf_text(&delta));
with_non_stream_total_timeout(plan, async move {
let key_upstream_metadata = read_windsurf_key_upstream_metadata(state, plan).await;
let prepared = prepare_windsurf_cascade(plan, input, key_upstream_metadata).await?;
let started_at = Instant::now();
let mut deltas = Vec::new();
let poll_result = poll_windsurf_cascade_with_transport_recovery(&prepared, |event| {
if let WindsurfPollEvent::TextDelta(delta) = event {
deltas.push(sanitize_windsurf_text(&delta));
}
Ok(())
})
.await?;
let elapsed_ms = started_at.elapsed().as_millis() as u64;
let content = deltas.concat();
let parsed_tool_calls = parse_and_filter_windsurf_tool_calls(&content, &prepared.input);
let mut tool_calls = poll_result.native_tool_calls;
tool_calls.extend(parsed_tool_calls.tool_calls);
let has_tool_calls = !tool_calls.is_empty();
let message = if has_tool_calls {
json!({
"role": "assistant",
"content": Value::Null,
"tool_calls": openai_tool_call_values(&tool_calls),
})
} else {
json!({
"role": "assistant",
"content": content,
})
};
let mut body_json = json!({
"id": format!("chatcmpl-{}", prepared.request_id),
"object": "chat.completion",
"created": current_unix_secs(),
"model": prepared.model,
"choices": [{
"index": 0,
"message": message,
"finish_reason": if has_tool_calls { "tool_calls" } else { "stop" },
}],
});
if let Some(usage) = poll_result.usage {
body_json["usage"] = windsurf_openai_usage_json(&usage);
}
Ok(())
})
.await?;
let elapsed_ms = started_at.elapsed().as_millis() as u64;
let content = deltas.concat();
let parsed_tool_calls = parse_and_filter_windsurf_tool_calls(&content, &prepared.input);
let mut tool_calls = poll_result.native_tool_calls;
tool_calls.extend(parsed_tool_calls.tool_calls);
let has_tool_calls = !tool_calls.is_empty();
let message = if has_tool_calls {
json!({
"role": "assistant",
"content": Value::Null,
"tool_calls": openai_tool_call_values(&tool_calls),
})
} else {
json!({
"role": "assistant",
"content": content,
})
};
let mut body_json = json!({
"id": format!("chatcmpl-{}", prepared.request_id),
"object": "chat.completion",
"created": current_unix_secs(),
"model": prepared.model,
"choices": [{
"index": 0,
"message": message,
"finish_reason": if has_tool_calls { "tool_calls" } else { "stop" },
}],
});
if let Some(usage) = poll_result.usage {
body_json["usage"] = windsurf_openai_usage_json(&usage);
}
Ok(Some(ExecutionResult {
request_id: prepared.request_id,
candidate_id: prepared.candidate_id,
status_code: 200,
headers: BTreeMap::from([("content-type".to_string(), "application/json".to_string())]),
body: Some(ResponseBody {
json_body: Some(body_json),
body_bytes_b64: None,
}),
telemetry: Some(ExecutionTelemetry {
ttfb_ms: None,
elapsed_ms: Some(elapsed_ms),
upstream_bytes: None,
}),
error: None,
}))
Ok(Some(ExecutionResult {
request_id: prepared.request_id,
candidate_id: prepared.candidate_id,
status_code: 200,
headers: BTreeMap::from([("content-type".to_string(), "application/json".to_string())]),
body: Some(ResponseBody {
json_body: Some(body_json),
body_bytes_b64: None,
}),
telemetry: Some(ExecutionTelemetry {
ttfb_ms: None,
elapsed_ms: Some(elapsed_ms),
upstream_bytes: None,
}),
error: None,
}))
})
.await
}
async fn prepare_windsurf_cascade(
@@ -1172,13 +1176,19 @@ async fn ensure_windsurf_language_server(
let mut command = Command::new(&binary_path);
command
.arg(format!("--api_server_url={}", codeium_api_url()))
.arg("--run_child")
.arg(format!("--server_port={port}"))
.arg(format!("--csrf_token={DEFAULT_CSRF_TOKEN}"))
.arg(format!("--register_user_url={DEFAULT_REGISTER_USER_URL}"))
.arg(format!("--codeium_dir={}", data_dir.display()))
.arg(format!("--database_dir={}", data_dir.join("db").display()))
.arg("--detect_proxy=false")
.env_clear()
.arg("--detect_proxy=false");
if !cfg!(target_os = "windows") {
command.env_clear();
}
command
.envs(language_server_env(proxy_url.as_deref()))
.stdin(Stdio::null())
.stdout(Stdio::null())
@@ -1191,7 +1201,8 @@ async fn ensure_windsurf_language_server(
))
})?;
if let Err(err) = wait_language_server_ready(port).await {
if let Err(err) = wait_language_server_ready(port, &mut child, stderr_log_path.as_deref()).await
{
let _ = child.kill();
return Err(err);
}
@@ -1379,7 +1390,7 @@ async fn windsurf_warmup_unary(
Err(err) if is_windsurf_cascade_transport_error(&err) => Err(err),
Err(err) => {
if stage == "UpdateWorkspaceTrust" {
error!(
warn!(
event_name = "windsurf_workspace_trust_update_failed",
log_type = "ops",
port,
@@ -1512,44 +1523,39 @@ async fn windsurf_grpc_unary(
) -> Result<Vec<u8>, ExecutionRuntimeTransportError> {
let url = format!("http://127.0.0.1:{port}{LS_SERVICE}/{method}");
let client = reqwest::Client::builder()
.http2_prior_knowledge()
.http1_only()
.timeout(timeout)
.build()
.map_err(ExecutionRuntimeTransportError::ClientBuild)?;
let response = client
.post(url)
.header("content-type", "application/grpc")
.header("te", "trailers")
.header("user-agent", "grpc-node/1.108.2")
.header("content-type", "application/proto")
.header("connect-protocol-version", "1")
.header("user-agent", "connect-es/1.5.0")
.header("x-codeium-csrf-token", csrf_token)
.body(grpc_frame(&payload))
.body(payload)
.send()
.await
.map_err(|err| {
ExecutionRuntimeTransportError::UpstreamRequest(format!(
"Windsurf gRPC {method} request failed: {}",
"Windsurf Connect {method} request failed: {}",
super::transport::format_upstream_request_error(&err)
))
})?;
let status = response.status();
let body = response.bytes().await.map_err(|err| {
ExecutionRuntimeTransportError::UpstreamRequest(format!(
"Windsurf gRPC {method} response read failed: {}",
"Windsurf Connect {method} response read failed: {}",
super::transport::format_upstream_request_error(&err)
))
})?;
if !status.is_success() {
return Err(ExecutionRuntimeTransportError::UpstreamRequest(format!(
"Windsurf gRPC {method} returned HTTP {status}: {}",
"Windsurf Connect {method} returned HTTP {status}: {}",
String::from_utf8_lossy(&body)
)));
}
let frames = extract_grpc_frames(&body);
if frames.is_empty() {
Ok(body.to_vec())
} else {
Ok(frames.concat())
}
Ok(body.to_vec())
}
fn detect_windsurf_request(
@@ -3897,10 +3903,27 @@ fn port_is_free(port: u16) -> bool {
TcpListener::bind(("127.0.0.1", port)).is_ok()
}
async fn wait_language_server_ready(port: u16) -> Result<(), ExecutionRuntimeTransportError> {
async fn wait_language_server_ready(
port: u16,
child: &mut Child,
stderr_log_path: Option<&Path>,
) -> Result<(), ExecutionRuntimeTransportError> {
let timeout = if cfg!(target_os = "windows") {
WINDOWS_LS_READY_TIMEOUT
} else {
LS_READY_TIMEOUT
};
let started = Instant::now();
let addr = SocketAddr::from(([127, 0, 0, 1], port));
while started.elapsed() < LS_READY_TIMEOUT {
while started.elapsed() < timeout {
if let Ok(Some(status)) = child.try_wait() {
let stderr_tail = stderr_log_path
.and_then(|path| read_log_tail(path, 8 * 1024))
.unwrap_or_else(|| "<stderr log unavailable>".to_string());
return Err(ExecutionRuntimeTransportError::UpstreamRequest(format!(
"Windsurf language server exited before port {port} became ready with status {status}; stderr tail: {stderr_tail}"
)));
}
if TcpStream::connect_timeout(&addr, Duration::from_millis(200)).is_ok() {
debug!(
event_name = "windsurf_language_server_port_ready",
@@ -3912,12 +3935,30 @@ async fn wait_language_server_ready(port: u16) -> Result<(), ExecutionRuntimeTra
}
tokio::time::sleep(Duration::from_millis(250)).await;
}
let child_status = child.try_wait().ok().flatten();
let stderr_tail = stderr_log_path
.and_then(|path| read_log_tail(path, 8 * 1024))
.unwrap_or_else(|| "<stderr log unavailable>".to_string());
Err(ExecutionRuntimeTransportError::UpstreamRequest(format!(
"Windsurf language server port {port} was not ready after {}ms",
LS_READY_TIMEOUT.as_millis()
"Windsurf language server port {port} was not ready after {}ms{}; stderr tail: {stderr_tail}",
timeout.as_millis(),
child_status
.map(|status| format!(" (child status: {status})"))
.unwrap_or_default()
)))
}
fn read_log_tail(path: &Path, max_bytes: usize) -> Option<String> {
let mut file = fs::File::open(path).ok()?;
let len = file.metadata().ok()?.len() as i64;
let max_bytes = max_bytes.max(1) as i64;
let start = len.saturating_sub(max_bytes);
file.seek(SeekFrom::Start(start as u64)).ok()?;
let mut buf = Vec::with_capacity((len - start) as usize);
file.read_to_end(&mut buf).ok()?;
Some(String::from_utf8_lossy(&buf).to_string())
}
fn language_server_data_dir(key: &str) -> PathBuf {
for env_key in ["WINDSURF_LS_DATA_DIR", "LS_DATA_DIR"] {
if let Some(path) = std::env::var_os(env_key).filter(|value| !value.is_empty()) {
@@ -3979,9 +4020,36 @@ fn language_server_env(proxy_url: Option<&str>) -> BTreeMap<String, String> {
}
}
}
if cfg!(target_os = "windows") {
for key in [
"USERPROFILE",
"APPDATA",
"LOCALAPPDATA",
"SystemRoot",
"WINDIR",
"ComSpec",
] {
if let Ok(value) = std::env::var(key) {
if !value.trim().is_empty() {
env.insert(key.to_string(), value);
}
}
}
}
if !env.contains_key("HOME") {
env.insert("HOME".to_string(), home_dir().display().to_string());
}
if cfg!(target_os = "windows") && !env.contains_key("USERPROFILE") {
if let Some(home) = std::env::var_os("USERPROFILE")
.or_else(|| std::env::var_os("HOME"))
.filter(|value| !value.is_empty())
{
env.insert(
"USERPROFILE".to_string(),
PathBuf::from(home).display().to_string(),
);
}
}
if let Some(proxy_url) = proxy_url.filter(|value| !value.trim().is_empty()) {
for key in ["HTTP_PROXY", "HTTPS_PROXY", "http_proxy", "https_proxy"] {
env.insert(key.to_string(), proxy_url.to_string());
@@ -3998,9 +4066,15 @@ fn codeium_api_url() -> String {
}
fn home_dir() -> PathBuf {
std::env::var_os("HOME")
.map(PathBuf::from)
.unwrap_or_else(|| PathBuf::from("."))
if let Some(home) = std::env::var_os("HOME").filter(|value| !value.is_empty()) {
return PathBuf::from(home);
}
if cfg!(target_os = "windows") {
if let Some(home) = std::env::var_os("USERPROFILE").filter(|value| !value.is_empty()) {
return PathBuf::from(home);
}
}
PathBuf::from(".")
}
fn first_string(body: &Value, keys: &[&str]) -> Option<String> {