Skip to content
Merged
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
6 changes: 3 additions & 3 deletions .repository-projection.json
Original file line number Diff line number Diff line change
Expand Up @@ -3,11 +3,11 @@
"projection": "deixic-code",
"projectionSchemaVersion": 1,
"sourceRepository": "dx-corp/mono",
"sourceSha": "ac81c8d7cff3a1002c8c5dbfac909bbaea58a854",
"sourceSha": "61ea3dea08e87af097277910b63a97fe6e216cf2",
"destinationRepository": "dx-corp/code",
"priorProjectedBase": "88f93e94824e24ef8cd97b3b9a3c4c02777b9549",
"priorProjectedBase": "11ca01744d12603c659827f5cec9e86e834b5da8",
"definitionDigest": "82936441c776e3e8edb5d215a75007ec9714a233f489d460075d79d5ef5ba32f",
"toolDigest": "c244d99199a7ae3eb8ff644a99462163c23b0bb6a83ef50af01efbdca0b81d04",
"contentDigest": "e57247ac09795028aad4cdee2023bcf4ad7a30fd7afbe126c5fad2e5bb70cb14",
"contentDigest": "9d215910206a5cea42c6e297cd0553a1ae73a53950fad655a696861809f80675",
"publicationEligible": true
}
71 changes: 3 additions & 68 deletions packages/ai-rs/src/openai.rs
Original file line number Diff line number Diff line change
Expand Up @@ -236,6 +236,9 @@ async fn send_with_response_open_timeout(

#[path = "openai/managed_gateway.rs"]
mod managed_gateway;
#[cfg(test)]
#[path = "openai/response_open_timeout_tests.rs"]
mod response_open_timeout_tests;
use managed_gateway::{
managed_gateway_error_retry_after, managed_gateway_receipt, managed_provider_tools_evidence,
};
Expand Down Expand Up @@ -6899,74 +6902,6 @@ data: {"type":"response.completed","response":{"output":[{"type":"message","cont
server.await.unwrap();
}

#[tokio::test]
async fn managed_gateway_response_open_timeout_stops_stalled_headers() {
use std::io::{Read, Write};
use std::net::TcpListener;

let listener = TcpListener::bind("127.0.0.1:0").expect("bind mock gateway");
let address = listener.local_addr().expect("mock gateway address");
std::thread::spawn(move || {
let (mut stream, _) = listener.accept().expect("accept gateway request");
let mut request = [0_u8; 4096];
let _ = stream.read(&mut request).expect("read gateway request");
std::thread::sleep(std::time::Duration::from_millis(150));
let _ = stream
.write_all(b"HTTP/1.1 200 OK\r\nContent-Length: 0\r\nConnection: close\r\n\r\n");
});

let request = reqwest::Client::new().get(format!("http://{address}/responses"));
let error = send_with_response_open_timeout(
request,
Some(std::time::Duration::from_millis(25)),
"managed gateway",
)
.await
.expect_err("stalled response opening must time out");

assert!(
error
.to_string()
.contains("managed gateway response headers timed out"),
"unexpected error: {error:#}"
);
}

#[tokio::test]
async fn managed_gateway_response_open_timeout_does_not_cover_stream_body() {
use std::io::{Read, Write};
use std::net::TcpListener;

let listener = TcpListener::bind("127.0.0.1:0").expect("bind mock gateway");
let address = listener.local_addr().expect("mock gateway address");
std::thread::spawn(move || {
let (mut stream, _) = listener.accept().expect("accept gateway request");
let mut request = [0_u8; 4096];
let _ = stream.read(&mut request).expect("read gateway request");
stream
.write_all(b"HTTP/1.1 200 OK\r\nContent-Length: 5\r\nConnection: close\r\n\r\n")
.expect("write response headers");
stream.flush().expect("flush response headers");
std::thread::sleep(std::time::Duration::from_millis(150));
stream.write_all(b"hello").expect("write delayed body");
});

let request = reqwest::Client::new().get(format!("http://{address}/responses"));
let response = send_with_response_open_timeout(
request,
Some(std::time::Duration::from_millis(50)),
"managed gateway",
)
.await
.expect("response headers should arrive within the open timeout");
let body = response
.text()
.await
.expect("body may continue past the response-open timeout");

assert_eq!(body, "hello");
}

#[test]
fn managed_gateway_scope_adds_workspace_header_without_changing_provider_reference() {
let client =
Expand Down
143 changes: 143 additions & 0 deletions packages/ai-rs/src/openai/response_open_timeout_tests.rs
Original file line number Diff line number Diff line change
@@ -0,0 +1,143 @@
//! Response-header deadlines preserve delayed-body streaming.
use super::send_with_response_open_timeout;

// Real loopback I/O must finish without Tokio advancing a paused clock to
// the open deadline just because the socket is briefly idle. A runnable
// task keeps time under the test's explicit control until it is aborted.
fn prevent_network_clock_auto_advance() -> tokio::task::JoinHandle<()> {
tokio::spawn(async {
loop {
tokio::task::yield_now().await;
}
})
}

#[tokio::test(start_paused = true)]
async fn managed_gateway_response_open_timeout_stops_stalled_headers() {
use tokio::io::{AsyncReadExt, AsyncWriteExt};

let clock = prevent_network_clock_auto_advance();
let listener = tokio::net::TcpListener::bind("127.0.0.1:0")
.await
.expect("bind mock gateway");
let address = listener.local_addr().expect("mock gateway address");
let (accepted, received) = tokio::sync::oneshot::channel();
let server = tokio::spawn(async move {
let (mut stream, _) = listener.accept().await.expect("accept gateway request");
let mut request = [0_u8; 4096];
assert!(
stream
.read(&mut request)
.await
.expect("read gateway request")
> 0
);
accepted.send(()).expect("request readiness");
tokio::time::sleep(std::time::Duration::from_millis(150)).await;
let _ = stream
.write_all(b"HTTP/1.1 200 OK\r\nContent-Length: 0\r\nConnection: close\r\n\r\n")
.await;
});

let request = reqwest::Client::builder()
.no_proxy()
.build()
.unwrap()
.get(format!("http://{address}/responses"));
let client = tokio::spawn(send_with_response_open_timeout(
request,
Some(std::time::Duration::from_millis(25)),
"managed gateway",
));
received.await.expect("request reached gateway");
assert!(
!client.is_finished(),
"headers are still pending before the deadline"
);
tokio::time::advance(std::time::Duration::from_millis(25)).await;
let error = client
.await
.unwrap()
.expect_err("stalled response opening must time out");

assert!(
error
.to_string()
.contains("managed gateway response headers timed out"),
"unexpected error: {error:#}"
);
tokio::time::advance(std::time::Duration::from_millis(125)).await;
server.await.expect("mock gateway server");
clock.abort();
}

#[tokio::test(start_paused = true)]
async fn managed_gateway_response_open_timeout_does_not_cover_stream_body() {
use tokio::io::{AsyncReadExt, AsyncWriteExt};

let clock = prevent_network_clock_auto_advance();
let listener = tokio::net::TcpListener::bind("127.0.0.1:0")
.await
.expect("bind mock gateway");
let address = listener.local_addr().expect("mock gateway address");
let (headers_written, received) = tokio::sync::oneshot::channel();
let server = tokio::spawn(async move {
let (mut stream, _) = listener.accept().await.expect("accept gateway request");
let mut request = [0_u8; 4096];
assert!(
stream
.read(&mut request)
.await
.expect("read gateway request")
> 0
);
stream
.write_all(b"HTTP/1.1 200 OK\r\nContent-Length: 5\r\nConnection: close\r\n\r\n")
.await
.expect("write response headers");
stream.flush().await.expect("flush response headers");
headers_written.send(()).expect("header readiness");
tokio::time::sleep(std::time::Duration::from_millis(150)).await;
stream
.write_all(b"hello")
.await
.expect("write delayed body");
});

let request = reqwest::Client::builder()
.no_proxy()
.build()
.unwrap()
.get(format!("http://{address}/responses"));
let client = tokio::spawn(send_with_response_open_timeout(
request,
Some(std::time::Duration::from_millis(50)),
"managed gateway",
));
received.await.expect("gateway wrote headers");
let response = client
.await
.unwrap()
.expect("response headers should arrive within the open timeout");
let body = response.text();
tokio::pin!(body);
assert!(futures::poll!(&mut body).is_pending());
tokio::time::advance(std::time::Duration::from_millis(50)).await;
assert!(
futures::poll!(&mut body).is_pending(),
"body must outlive the 50 ms open timeout"
);
tokio::time::advance(std::time::Duration::from_millis(99)).await;
assert!(
futures::poll!(&mut body).is_pending(),
"body remains pending until its 150 ms delay"
);
tokio::time::advance(std::time::Duration::from_millis(1)).await;
let body = body
.await
.expect("body may continue past the response-open timeout");

assert_eq!(body, "hello");
server.await.expect("mock gateway server");
clock.abort();
}
17 changes: 11 additions & 6 deletions packages/local-host-rs/src/native_credentials.rs
Original file line number Diff line number Diff line change
@@ -1,14 +1,17 @@
//! Native credential access shared by OAuth, connections, and MCP.
//!
//! Unit tests cannot open the developer's credential store. Process-based
//! fixtures use MAESTRO_DISABLE_KEYCHAIN=1 to enforce the same boundary.
//! Unit tests and opt-in test-support dependency builds cannot open the
//! developer's credential store. Other process-based fixtures use
//! MAESTRO_DISABLE_KEYCHAIN=1 to enforce the same boundary.
use anyhow::{Result, bail};

pub(crate) fn entry(service: &str, account: &str) -> Result<keyring::Entry> {
let disabled = std::env::var("MAESTRO_DISABLE_KEYCHAIN").ok();
open_with_policy(cfg!(test), disabled.as_deref(), || {
keyring::Entry::new(service, account).map_err(Into::into)
})
open_with_policy(
cfg!(any(test, feature = "test-support")),
disabled.as_deref(),
|| keyring::Entry::new(service, account).map_err(Into::into),
)
}

fn open_with_policy<T>(
Expand All @@ -17,7 +20,9 @@ fn open_with_policy<T>(
open: impl FnOnce() -> Result<T>,
) -> Result<T> {
if unit_test {
bail!("native credential access is disabled in unit tests; inject a test secret backend");
bail!(
"native credential access is disabled in unit tests and test-support builds; inject a test secret backend"
);
}
if disabled == Some("1") {
bail!("native credential access is disabled by MAESTRO_DISABLE_KEYCHAIN=1");
Expand Down
50 changes: 50 additions & 0 deletions packages/local-host-rs/tests/embedding.rs
Original file line number Diff line number Diff line change
Expand Up @@ -12,6 +12,56 @@ use maestro_local_host::embedding::{
};
use maestro_local_host::state::ApprovalMode;

#[test]
fn test_support_dependency_cannot_open_native_credentials() {
const CHILD: &str = "MAESTRO_TEST_NATIVE_CREDENTIAL_BOUNDARY";
if std::env::var_os(CHILD).is_some() {
// This integration binary links local-host without cfg(test), just as
// TUI tests do. Force the real public storage path without the runtime
// disable flag so the compiled test-support boundary must reject it.
let error = match maestro_local_host::init_cli::load_evalops_snapshot() {
Err(error) => error,
Ok(_) => panic!("test-support dependency opened native credential storage"),
};
assert!(format!("{error:#}").contains("disabled in unit tests and test-support builds"));
return;
}
let home = tempfile::tempdir().expect("isolated credential home");
let mut child = std::process::Command::new(std::env::current_exe().unwrap())
.args([
"--exact",
"test_support_dependency_cannot_open_native_credentials",
"--nocapture",
])
.env(CHILD, "1")
.env("MAESTRO_HOME", home.path())
.env("MAESTRO_OAUTH_STORAGE_MODE", "keychain")
.env_remove("MAESTRO_DISABLE_KEYCHAIN")
.spawn()
.expect("spawn isolated credential boundary fixture");
let deadline = std::time::Instant::now() + Duration::from_secs(5);
loop {
if let Some(status) = child
.try_wait()
.expect("wait for credential boundary fixture")
{
assert!(
status.success(),
"credential boundary fixture failed: {status}"
);
break;
}
if std::time::Instant::now() >= deadline {
child
.kill()
.expect("stop stalled credential boundary fixture");
child.wait().expect("reap credential boundary fixture");
panic!("credential boundary fixture did not reject native storage promptly");
}
std::thread::sleep(Duration::from_millis(10));
}
}

async fn next_event_matching(
session: &mut EmbeddedAgentSession,
predicate: impl Fn(&FromAgent) -> bool,
Expand Down
Loading
Loading