diff --git a/crates/xrpc/tests/search.rs b/crates/xrpc/tests/search.rs new file mode 100644 index 0000000..5916f50 --- /dev/null +++ b/crates/xrpc/tests/search.rs @@ -0,0 +1,478 @@ +use std::sync::Arc; + +use axum::body::{Body, to_bytes}; +use bobbin_edge_index::{Coverage, CoverageWatch, EdgeStore, HydrantCursor}; +use bobbin_knot_proxy::{KnotProxy, KnotProxyConfig}; +use bobbin_record_lru::{CacheCapacity, LruRecordStore}; +use bobbin_search::{DEFAULT_WRITER_HEAP_BYTES, SearchIndex}; +use bobbin_slingshot_client::SlingshotClient; +use bobbin_types::search::{SearchDoc, SearchSink}; +use bobbin_xrpc::{AppState, router}; +use http::{Request, StatusCode}; +use jacquard_common::DefaultStr; +use jacquard_common::types::nsid::Nsid; +use jacquard_common::types::string::AtUri; +use serde_json::{Value, json}; +use tower::ServiceExt; +use url::Url; +use url::form_urlencoded::byte_serialize; +use wiremock::matchers::{method, path, query_param}; +use wiremock::{Mock, MockServer, ResponseTemplate}; + +const CID: &str = "bafyreieqygohnz2zqyvtvktbjpvhutphobcmbsnt4q5lc36ri7vpcmoz4i"; + +fn at(s: &str) -> AtUri { + AtUri::new_owned(s).unwrap() +} + +fn nsid(s: &'static str) -> Nsid { + Nsid::new_static(s).unwrap() +} + +fn enc(s: &str) -> String { + byte_serialize(s.as_bytes()).collect() +} + +struct Harness { + server: MockServer, + coverage: Arc, + search: Arc, + state: AppState, +} + +impl Harness { + async fn new() -> Self { + let server = MockServer::start().await; + let coverage = Arc::new(CoverageWatch::new()); + let search = Arc::new(SearchIndex::new(DEFAULT_WRITER_HEAP_BYTES).unwrap()); + let state = AppState::new( + Arc::new(LruRecordStore::new(CacheCapacity::from_bytes(64 * 1024))), + SlingshotClient::new(Url::parse(&server.uri()).unwrap()).unwrap(), + Arc::new(EdgeStore::new()), + coverage.clone(), + Arc::new(KnotProxy::new(KnotProxyConfig::default()).unwrap()), + search.clone(), + ); + Self { + server, + coverage, + search, + state, + } + } + + async fn index_issue(&self, did: &str, rkey: &str, title: &str, body: &str) { + let uri = format!("at://{did}/sh.tangled.repo.issue/{rkey}"); + self.search + .upsert(SearchDoc { + uri: at(&uri), + nsid: nsid("sh.tangled.repo.issue"), + title: title.to_owned(), + body: body.to_owned(), + }) + .await; + self.search.flush().await; + self.mount_issue(did, rkey, title, body).await; + } + + async fn index_string(&self, did: &str, rkey: &str, filename: &str, contents: &str) { + let uri = format!("at://{did}/sh.tangled.string/{rkey}"); + self.search + .upsert(SearchDoc { + uri: at(&uri), + nsid: nsid("sh.tangled.string"), + title: filename.to_owned(), + body: contents.to_owned(), + }) + .await; + self.search.flush().await; + self.mount_string(did, rkey, filename, contents).await; + } + + async fn mount_issue(&self, did: &str, rkey: &str, title: &str, body: &str) { + let uri = format!("at://{did}/sh.tangled.repo.issue/{rkey}"); + let value = json!({ + "$type": "sh.tangled.repo.issue", + "repoDid": "did:plc:abalone", + "title": title, + "body": body, + "createdAt": "2026-05-01T00:00:00Z" + }); + Mock::given(method("GET")) + .and(path("/xrpc/com.atproto.repo.getRecord")) + .and(query_param("repo", did)) + .and(query_param("collection", "sh.tangled.repo.issue")) + .and(query_param("rkey", rkey)) + .respond_with(ResponseTemplate::new(200).set_body_json(json!({ + "uri": uri, + "cid": CID, + "value": value, + }))) + .mount(&self.server) + .await; + } + + async fn mount_string(&self, did: &str, rkey: &str, filename: &str, contents: &str) { + let uri = format!("at://{did}/sh.tangled.string/{rkey}"); + let value = json!({ + "$type": "sh.tangled.string", + "filename": filename, + "description": "field notes", + "contents": contents, + "createdAt": "2026-05-01T00:00:00Z" + }); + Mock::given(method("GET")) + .and(path("/xrpc/com.atproto.repo.getRecord")) + .and(query_param("repo", did)) + .and(query_param("collection", "sh.tangled.string")) + .and(query_param("rkey", rkey)) + .respond_with(ResponseTemplate::new(200).set_body_json(json!({ + "uri": uri, + "cid": CID, + "value": value, + }))) + .mount(&self.server) + .await; + } + + fn warming(&self, events: u64, cursor: u64) { + self.coverage.update(|_| Coverage::Warming { + events_processed: events, + last_cursor: HydrantCursor::new(cursor), + }); + } + + fn promote_ready(&self, events: u64, cursor: u64) { + self.coverage.update(|_| Coverage::Ready { + events_processed: events, + last_cursor: HydrantCursor::new(cursor), + }); + } +} + +fn search_request(extras: &[(&str, &str)]) -> Request { + let qs = extras + .iter() + .map(|(k, v)| format!("{k}={}", enc(v))) + .collect::>() + .join("&"); + Request::builder() + .uri(format!("/xrpc/sh.tangled.search.query?{qs}")) + .body(Body::empty()) + .unwrap() +} + +async fn json_response(resp: axum::response::Response) -> (StatusCode, Value) { + let status = resp.status(); + let bytes = to_bytes(resp.into_body(), 1 << 20).await.unwrap(); + let parsed: Value = serde_json::from_slice(&bytes).expect("JSON body"); + (status, parsed) +} + +#[tokio::test] +async fn empty_query_returns_no_hits_with_warming_envelope() { + let h = Harness::new().await; + let app = router(h.state.clone()); + let resp = app + .oneshot(search_request(&[("q", "barnacle")])) + .await + .unwrap(); + let (status, body) = json_response(resp).await; + assert_eq!(status, StatusCode::OK); + assert_eq!(body["hits"], json!([])); + assert_eq!(body["coverage"]["ready"], json!(false)); + assert_eq!(body["coverage"]["eventsProcessed"], json!(0)); + assert!(body["cursor"].is_null()); +} + +#[tokio::test] +async fn indexed_issue_hydrates_typed_value() { + let h = Harness::new().await; + h.warming(7, 42); + h.index_issue( + "did:plc:nel", + "abcabcabcabcz", + "barnacle pagination overflow", + "scrolling resets when the cursor wraps", + ) + .await; + let app = router(h.state.clone()); + let resp = app + .oneshot(search_request(&[("q", "barnacle")])) + .await + .unwrap(); + let (status, body) = json_response(resp).await; + assert_eq!(status, StatusCode::OK); + let hits = body["hits"].as_array().expect("hits array"); + assert_eq!(hits.len(), 1); + assert_eq!( + hits[0]["uri"].as_str().unwrap(), + "at://did:plc:nel/sh.tangled.repo.issue/abcabcabcabcz", + ); + assert_eq!(hits[0]["nsid"], json!("sh.tangled.repo.issue")); + assert_eq!(hits[0]["cid"], CID); + assert_eq!( + hits[0]["value"]["$type"], + json!("sh.tangled.repo.issue"), + "hit value is typed via SearchableRecord serialization", + ); + assert_eq!( + hits[0]["value"]["title"], + json!("barnacle pagination overflow"), + ); + assert_eq!(hits[0]["value"]["repoDid"], json!("did:plc:abalone")); + assert!(hits[0]["score"].as_f64().unwrap() > 0.0); + assert_eq!(body["coverage"]["ready"], json!(false)); + assert_eq!(body["coverage"]["eventsProcessed"], json!(7)); + assert_eq!(body["coverage"]["lastCursor"], json!(42)); +} + +#[tokio::test] +async fn nsid_filter_narrows_to_single_collection() { + let h = Harness::new().await; + h.index_issue("did:plc:nel", "i1", "anemone tide", "high water") + .await; + h.index_string("did:plc:teq", "k1", "anemone.md", "anemone notes") + .await; + let app = router(h.state.clone()); + let resp = app + .clone() + .oneshot(search_request(&[ + ("q", "anemone"), + ("nsid", "sh.tangled.string"), + ])) + .await + .unwrap(); + let (status, body) = json_response(resp).await; + assert_eq!(status, StatusCode::OK); + let hits = body["hits"].as_array().unwrap(); + assert_eq!(hits.len(), 1); + assert_eq!(hits[0]["nsid"], json!("sh.tangled.string")); + assert_eq!(hits[0]["value"]["$type"], json!("sh.tangled.string")); + assert_eq!(hits[0]["value"]["filename"], json!("anemone.md")); +} + +#[tokio::test] +async fn coverage_promotion_does_not_change_hit_set() { + let h = Harness::new().await; + h.index_issue("did:plc:nel", "i1", "limpet survey", "tidal pools") + .await; + let app = router(h.state.clone()); + h.warming(3, 5); + let (_, before) = json_response( + app.clone() + .oneshot(search_request(&[("q", "limpet")])) + .await + .unwrap(), + ) + .await; + assert_eq!(before["coverage"]["ready"], json!(false)); + let before_hits = before["hits"].as_array().unwrap().clone(); + assert_eq!(before_hits.len(), 1); + + h.promote_ready(9, 11); + let (_, after) = json_response( + app.oneshot(search_request(&[("q", "limpet")])) + .await + .unwrap(), + ) + .await; + assert_eq!(after["coverage"]["ready"], json!(true)); + assert_eq!(after["coverage"]["lastCursor"], json!(11)); + let after_hits = after["hits"].as_array().unwrap(); + assert_eq!(after_hits.len(), before_hits.len()); + assert_eq!(after_hits[0]["uri"], before_hits[0]["uri"]); +} + +#[tokio::test] +async fn pagination_round_trips_via_cursor() { + let h = Harness::new().await; + let names = ["nel", "olaren", "teq", "lyna", "bailey"]; + for (i, owner) in names.iter().enumerate() { + h.index_issue( + &format!("did:plc:{owner}"), + &format!("r{i}"), + "anemone tides", + "shell", + ) + .await; + } + let app = router(h.state.clone()); + let (_, page1) = json_response( + app.clone() + .oneshot(search_request(&[("q", "anemone"), ("limit", "2")])) + .await + .unwrap(), + ) + .await; + assert_eq!(page1["hits"].as_array().unwrap().len(), 2); + let cursor = page1["cursor"].as_str().expect("more pages").to_owned(); + + let (_, page2) = json_response( + app.clone() + .oneshot(search_request(&[ + ("q", "anemone"), + ("limit", "2"), + ("cursor", &cursor), + ])) + .await + .unwrap(), + ) + .await; + assert_eq!(page2["hits"].as_array().unwrap().len(), 2); + let cursor2 = page2["cursor"].as_str().expect("more pages").to_owned(); + + let (_, page3) = json_response( + app.oneshot(search_request(&[ + ("q", "anemone"), + ("limit", "2"), + ("cursor", &cursor2), + ])) + .await + .unwrap(), + ) + .await; + assert_eq!(page3["hits"].as_array().unwrap().len(), 1); + assert!(page3["cursor"].is_null()); +} + +#[tokio::test] +async fn empty_q_returns_400() { + let h = Harness::new().await; + let app = router(h.state.clone()); + let resp = app + .oneshot(search_request(&[("q", " ")])) + .await + .unwrap(); + let (status, body) = json_response(resp).await; + assert_eq!(status, StatusCode::BAD_REQUEST); + assert_eq!(body["error"], json!("InvalidRequest")); +} + +#[tokio::test] +async fn invalid_cursor_returns_400() { + let h = Harness::new().await; + let app = router(h.state.clone()); + let resp = app + .oneshot(search_request(&[("q", "anything"), ("cursor", "not-hex")])) + .await + .unwrap(); + let (status, body) = json_response(resp).await; + assert_eq!(status, StatusCode::BAD_REQUEST); + assert_eq!(body["error"], json!("InvalidRequest")); +} + +#[tokio::test] +async fn invalid_nsid_returns_400() { + let h = Harness::new().await; + let app = router(h.state.clone()); + let resp = app + .oneshot(search_request(&[ + ("q", "anything"), + ("nsid", "not a real nsid"), + ])) + .await + .unwrap(); + let (status, body) = json_response(resp).await; + assert_eq!(status, StatusCode::BAD_REQUEST); + assert_eq!(body["error"], json!("InvalidRequest")); +} + +#[tokio::test] +async fn tombstoned_hit_silently_dropped_from_results() { + let h = Harness::new().await; + h.search + .upsert(SearchDoc { + uri: at("at://did:plc:nel/sh.tangled.repo.issue/i1"), + nsid: nsid("sh.tangled.repo.issue"), + title: "kelp".to_owned(), + body: "ocean".to_owned(), + }) + .await; + h.search.flush().await; + h.index_issue("did:plc:teq", "i2", "kelp survives", "still here") + .await; + let app = router(h.state.clone()); + let resp = app + .oneshot(search_request(&[("q", "kelp")])) + .await + .unwrap(); + let (status, body) = json_response(resp).await; + assert_eq!(status, StatusCode::OK); + let hits = body["hits"].as_array().expect("hits array"); + assert_eq!(hits.len(), 1, "tombstoned hit dropped, sibling kept"); + assert_eq!( + hits[0]["uri"].as_str().unwrap(), + "at://did:plc:teq/sh.tangled.repo.issue/i2", + ); +} + +#[tokio::test] +async fn upstream_5xx_during_hydration_propagates_as_502() { + let h = Harness::new().await; + h.search + .upsert(SearchDoc { + uri: at("at://did:plc:nel/sh.tangled.repo.issue/i1"), + nsid: nsid("sh.tangled.repo.issue"), + title: "abalone".to_owned(), + body: "shell".to_owned(), + }) + .await; + h.search.flush().await; + Mock::given(method("GET")) + .and(path("/xrpc/com.atproto.repo.getRecord")) + .and(query_param("repo", "did:plc:nel")) + .and(query_param("collection", "sh.tangled.repo.issue")) + .and(query_param("rkey", "i1")) + .respond_with(ResponseTemplate::new(503)) + .mount(&h.server) + .await; + let app = router(h.state.clone()); + let resp = app + .oneshot(search_request(&[("q", "abalone")])) + .await + .unwrap(); + let (status, body) = json_response(resp).await; + assert_eq!(status, StatusCode::BAD_GATEWAY); + assert_eq!(body["error"], json!("UpstreamFailed")); +} + +#[tokio::test] +async fn second_query_short_circuits_via_lru_without_re_querying_slingshot() { + let h = Harness::new().await; + h.search + .upsert(SearchDoc { + uri: at("at://did:plc:nel/sh.tangled.repo.issue/i1"), + nsid: nsid("sh.tangled.repo.issue"), + title: "kelp".to_owned(), + body: "ocean".to_owned(), + }) + .await; + h.search.flush().await; + Mock::given(method("GET")) + .and(path("/xrpc/com.atproto.repo.getRecord")) + .and(query_param("repo", "did:plc:nel")) + .and(query_param("collection", "sh.tangled.repo.issue")) + .and(query_param("rkey", "i1")) + .respond_with(ResponseTemplate::new(200).set_body_json(json!({ + "uri": "at://did:plc:nel/sh.tangled.repo.issue/i1", + "cid": CID, + "value": { + "$type": "sh.tangled.repo.issue", + "repoDid": "did:plc:abalone", + "title": "kelp", + "body": "ocean", + "createdAt": "2026-05-01T00:00:00Z" + } + }))) + .expect(1) + .mount(&h.server) + .await; + let app = router(h.state.clone()); + let _ = app + .clone() + .oneshot(search_request(&[("q", "kelp")])) + .await + .unwrap(); + let _ = app.oneshot(search_request(&[("q", "kelp")])).await.unwrap(); +}