//! Tests for the real HTTP-backed embedder (LOT C1a, §14.5.3), gated by the //! `vector-http` feature. They exercise [`HttpEmbedder`] against a **minimal, //! one-shot, in-process HTTP server** (raw tokio `TcpListener`, no new test //! dependency) and the [`detect_ollama`] probe. //! //! The whole file is compiled out unless `--features vector-http` is set, so the //! default dependency-free build is unaffected. #![cfg(feature = "vector-http")] use tokio::io::{AsyncReadExt, AsyncWriteExt}; use tokio::net::TcpListener; use domain::ports::{Embedder, EmbedderError}; use domain::profile::{EmbedderProfile, EmbedderStrategy}; use infrastructure::{detect_ollama, HttpEmbedder}; /// Spawns a one-shot HTTP server on `127.0.0.1:0` that, for the next single /// connection, reads the full request (honouring `Content-Length`) then writes /// `response` verbatim and closes. Returns the bound `base` URL (`http://host:port`). async fn one_shot_server(response: &'static str) -> String { let listener = TcpListener::bind("127.0.0.1:0").await.unwrap(); let addr = listener.local_addr().unwrap(); tokio::spawn(async move { if let Ok((mut stream, _)) = listener.accept().await { drain_request(&mut stream).await; let _ = stream.write_all(response.as_bytes()).await; let _ = stream.flush().await; } }); format!("http://{addr}") } /// Reads an HTTP request off `stream` until headers (and any `Content-Length` /// body) are fully consumed, so the client never sees a premature reset. async fn drain_request(stream: &mut tokio::net::TcpStream) { let mut buf = Vec::new(); let mut tmp = [0u8; 1024]; loop { let headers_end = find_subslice(&buf, b"\r\n\r\n"); if let Some(h) = headers_end { let header_text = String::from_utf8_lossy(&buf[..h]).to_ascii_lowercase(); let content_len = header_text .lines() .find_map(|l| l.strip_prefix("content-length:")) .and_then(|v| v.trim().parse::().ok()) .unwrap_or(0); if buf.len() >= h + 4 + content_len { return; } } match stream.read(&mut tmp).await { Ok(0) => return, Ok(n) => buf.extend_from_slice(&tmp[..n]), Err(_) => return, } } } fn find_subslice(haystack: &[u8], needle: &[u8]) -> Option { haystack.windows(needle.len()).position(|w| w == needle) } /// Builds an HTTP `200 OK` response with a JSON body and the right `Content-Length`. fn ok_json(body: &str) -> String { format!( "HTTP/1.1 200 OK\r\nContent-Type: application/json\r\nContent-Length: {}\r\n\r\n{body}", body.len() ) } fn local_server_profile(endpoint: &str, dimension: usize) -> EmbedderProfile { EmbedderProfile::new( "test-local", "Test Local", EmbedderStrategy::LocalServer, Some("test-model".to_string()), Some(endpoint.to_string()), None, dimension, ) .unwrap() } #[tokio::test] async fn http_embedder_parses_vectors_and_restores_input_order() { // The server returns the two embeddings with their `index` fields swapped; the // embedder must sort by `index` so the output matches the *input* order. let body = r#"{"data":[ {"embedding":[0.0,1.0],"index":1}, {"embedding":[1.0,0.0],"index":0} ]}"#; let base = one_shot_server_leaked(ok_json(body)).await; let embedder = HttpEmbedder::from_profile(&local_server_profile(&base, 2)); let out = embedder .embed(&["first".to_string(), "second".to_string()]) .await .expect("embed must succeed"); assert_eq!( out, vec![vec![1.0, 0.0], vec![0.0, 1.0]], "input order restored by index" ); } #[tokio::test] async fn http_embedder_empty_input_short_circuits_without_network() { // A closed/never-bound endpoint: no request must be made for empty input. let embedder = HttpEmbedder::from_profile(&local_server_profile("http://127.0.0.1:1/v1/embeddings", 4)); let out = embedder .embed(&[]) .await .expect("empty input ⇒ empty output, no I/O"); assert!(out.is_empty()); } #[tokio::test] async fn http_embedder_non_2xx_is_unavailable() { let resp = "HTTP/1.1 500 Internal Server Error\r\nContent-Length: 0\r\n\r\n".to_string(); let base = one_shot_server_leaked(resp).await; let embedder = HttpEmbedder::from_profile(&local_server_profile(&base, 2)); let err = embedder.embed(&["x".to_string()]).await.unwrap_err(); assert!(matches!(err, EmbedderError::Unavailable(_)), "got {err:?}"); } #[tokio::test] async fn http_embedder_unreachable_host_is_unavailable() { // Port 1 on loopback: nothing listens ⇒ connection refused ⇒ Unavailable. let embedder = HttpEmbedder::from_profile(&local_server_profile("http://127.0.0.1:1/v1/embeddings", 2)); let err = embedder.embed(&["x".to_string()]).await.unwrap_err(); assert!(matches!(err, EmbedderError::Unavailable(_)), "got {err:?}"); } #[tokio::test] async fn http_embedder_malformed_body_is_io() { let base = one_shot_server_leaked(ok_json("not json at all")).await; let embedder = HttpEmbedder::from_profile(&local_server_profile(&base, 2)); let err = embedder.embed(&["x".to_string()]).await.unwrap_err(); assert!(matches!(err, EmbedderError::Io(_)), "got {err:?}"); } #[tokio::test] async fn http_embedder_count_mismatch_is_io() { // Two inputs but the server returns a single embedding. let body = r#"{"data":[{"embedding":[1.0,0.0],"index":0}]}"#; let base = one_shot_server_leaked(ok_json(body)).await; let embedder = HttpEmbedder::from_profile(&local_server_profile(&base, 2)); let err = embedder .embed(&["a".to_string(), "b".to_string()]) .await .unwrap_err(); assert!(matches!(err, EmbedderError::Io(_)), "got {err:?}"); } #[tokio::test] async fn http_embedder_dimension_mismatch_is_io() { // Profile declares dimension 4 but the server returns a length-2 vector. let body = r#"{"data":[{"embedding":[1.0,0.0],"index":0}]}"#; let base = one_shot_server_leaked(ok_json(body)).await; let embedder = HttpEmbedder::from_profile(&local_server_profile(&base, 4)); let err = embedder.embed(&["x".to_string()]).await.unwrap_err(); assert!(matches!(err, EmbedderError::Io(_)), "got {err:?}"); } #[tokio::test] async fn http_embedder_api_strategy_missing_key_is_unavailable() { // `api` strategy whose configured key env var is guaranteed unset ⇒ Unavailable // *before* any network call (we point at a dead endpoint to prove no request). let profile = EmbedderProfile::new( "test-api", "Test API", EmbedderStrategy::Api, Some("model".to_string()), Some("http://127.0.0.1:1/v1/embeddings".to_string()), Some("IDEA_TEST_DEFINITELY_UNSET_KEY_VAR".to_string()), 2, ) .unwrap(); let embedder = HttpEmbedder::from_profile(&profile); let err = embedder.embed(&["x".to_string()]).await.unwrap_err(); assert!(matches!(err, EmbedderError::Unavailable(_)), "got {err:?}"); } #[tokio::test] async fn detect_ollama_true_when_tags_endpoint_ok() { let base = one_shot_server_leaked(ok_json(r#"{"models":[]}"#)).await; assert!(detect_ollama(&base).await, "a 200 on /api/tags ⇒ detected"); } #[tokio::test] async fn detect_ollama_false_when_nothing_listening() { // Port 1 on loopback: connection refused ⇒ not detected, never panics. assert!(!detect_ollama("http://127.0.0.1:1").await); } /// Like [`one_shot_server`] but takes an owned `String` and leaks it to obtain the /// `'static` lifetime the spawned task needs (test-only; the process is short-lived). async fn one_shot_server_leaked(response: String) -> String { let leaked: &'static str = Box::leak(response.into_boxed_str()); one_shot_server(leaked).await }