use std::collections::{HashMap, HashSet}; use axum::extract::State; use bobbin_codesearch::{CodeSearchError, FileMatch, Range, extract_did}; use futures::{StreamExt as _, stream}; use jacquard_axum::service_auth::ExtractServiceAuth; use jacquard_axum::{ExtractXrpc, XrpcResponse}; use jacquard_common::IntoStatic as _; use jacquard_common::types::did::Did; use lexicons::org_tangled::temp::search::search_code::{ self, Chunk, FileResult, Highlight, SearchCodeOutput, }; use crate::{AppState, FETCH_CONCURRENCY, XrpcError, view::build_repo_view_basic}; const DEFAULT_LIMIT: i64 = 50; const MAX_LIMIT: i64 = 100; pub(crate) async fn search_code( State(state): State, ExtractServiceAuth(auth): ExtractServiceAuth, ExtractXrpc(params): ExtractXrpc, ) -> Result, XrpcError> { let viewer = auth.did().into_static(); search_code_inner(&state, Some(&viewer), params).await } async fn search_code_inner( state: &AppState, viewer: Option<&Did>, params: search_code::SearchCode, ) -> Result, XrpcError> { let Some(codesearch) = state.codesearch.as_ref() else { return Err(XrpcError::NotImplemented("code search is not configured")); }; let q = params.q.trim(); if q.is_empty() { return Err(XrpcError::InvalidParams("q must not be empty".into())); } // scope by the repo's own did (meta.did) let query = match ¶ms.repo { Some(did) => format!("meta.did:{did} {q}"), None => q.to_owned(), }; let limit = params.limit.unwrap_or(DEFAULT_LIMIT).clamp(1, MAX_LIMIT) as usize; let offset = params .cursor .as_deref() .map(str::parse::) .transpose() .map_err(|e| XrpcError::InvalidParams(format!("cursor: {e}")))? .unwrap_or(0); let _permit = state.heavy_permit()?; let page = codesearch .search(&query, offset, limit) .await .map_err(map_codesearch_err)?; let mut matches: Vec<(Did, FileMatch)> = page .files .into_iter() .filter_map(|file| { let did = page .repo_urls .get(&file.repository) .and_then(|template| extract_did(template))?; Some((did, file)) }) .collect(); if let Some(wanted) = ¶ms.repo { matches.retain(|(did, _)| did == wanted); } let views = hydrate_repos(state, viewer, &matches).await; let mut results = Vec::with_capacity(matches.len()); for (did, file) in &matches { let Some(repo) = views.get(did).cloned() else { continue; }; let mut ranges = None; let mut chunks = Vec::new(); for chunk in &file.chunk_matches { if chunk.file_name { ranges = Some(chunk_highlights(&file.file_name, 1, &chunk.ranges)); break; } let line_start = chunk.content_start.line_number; chunks.push(Chunk { highlights: Some(chunk_highlights( &chunk.content, line_start.into(), &chunk.ranges, )), content: chunk.content.clone().into(), line_start: line_start.into(), extra_data: Default::default(), }); } results.push(FileResult { repo, path: file.file_name.clone().into(), commit: file.version.clone().into(), branches: (!file.branches.is_empty()) .then(|| file.branches.iter().map(|b| b.as_str().into()).collect()), language: (!file.language.is_empty()).then(|| file.language.as_str().into()), ranges, chunks: (!chunks.is_empty()).then_some(chunks), extra_data: Default::default(), }); } Ok(XrpcResponse(SearchCodeOutput { cursor: page .has_more .then(|| (offset.saturating_add(limit)).to_string().into()), results, extra_data: Default::default(), })) } async fn hydrate_repos( state: &AppState, viewer: Option<&Did>, matches: &[(Did, FileMatch)], ) -> HashMap { let mut seen = HashSet::new(); let dids: Vec = matches .iter() .filter(|(did, _)| seen.insert(did.clone())) .map(|(did, _)| did.clone()) .collect(); stream::iter(dids.into_iter().map(|did| async move { match build_repo_view_basic(state, viewer, &did).await { Ok(view) => Some((did, view)), Err(err) => { tracing::warn!(%did, error = %err, "skipping code search result, repo hydration failed"); None } } })) .buffered(FETCH_CONCURRENCY) .filter_map(std::future::ready) .collect() .await } /// Maps zoekt's 1-based (line, rune-column) ranges onto byte offsets within the chunk's content. fn chunk_highlights(content: &str, start_line: i64, ranges: &[Range]) -> Vec { let mut offset = 0usize; let lines: Vec<(usize, &str)> = content .split_inclusive('\n') .map(|line| { let start = offset; offset += line.len(); (start, line) }) .collect(); let start_line = start_line.max(1); // maps a 1-based (line, rune column) to a byte offset in content let byte_at = |line: u32, column: u32| -> Option { let idx = usize::try_from(i64::from(line) - start_line).ok()?; let (offset, text) = *lines.get(idx)?; let col = column.saturating_sub(1) as usize; // a column past the end of the line clamps to just before its newline Some( offset + text .char_indices() .nth(col) .map_or_else(|| text.trim_end_matches('\n').len(), |(byte, _)| byte), ) }; ranges .iter() .filter_map(|range| { let start = byte_at(range.start.line_number, range.start.column)?; let end = byte_at(range.end.line_number, range.end.column)?; (start < end).then(|| Highlight { start: start as i64, end: end as i64, extra_data: Default::default(), }) }) .collect() } fn map_codesearch_err(err: CodeSearchError) -> XrpcError { use CodeSearchError as E; match err { E::BadQuery(e) => XrpcError::InvalidParams(format!("query: {e}")), e @ (E::Network(_) | E::Upstream { .. }) => { XrpcError::UpstreamUnavailable(format!("code search: {e}")) } e @ E::Decode(_) => XrpcError::InvalidRecord(format!("code search: {e}")), e @ (E::BadScheme(_) | E::Build(_)) => XrpcError::Internal(format!("code search: {e}")), } } #[cfg(test)] mod tests { use super::*; use bobbin_codesearch::{CodeSearch, Location}; use bobbin_edge_index::{CoverageWatch, EdgeStore, StateIndex}; use bobbin_knot_proxy::{KnotHttpConfig, KnotProxy, KnotProxyConfig}; use bobbin_record_lru::{CacheCapacity, LruRecordStore}; use bobbin_resolver::RepoIdResolver; use bobbin_runtime::{RuntimeHasher, SystemClock}; use bobbin_search::{DEFAULT_WRITER_HEAP_BYTES, SearchIndex, SearchReader}; use bobbin_slingshot_client::SlingshotClient; use jacquard_common::types::datetime::Datetime; use jacquard_common::types::recordkey::Rkey; use serde_json::{Value, json}; use std::sync::Arc; use url::Url; use wiremock::matchers::{method, path, query_param}; use wiremock::{Mock, MockServer, ResponseTemplate}; const CID: &str = "bafyreieqygohnz2zqyvtvktbjpvhutphobcmbsnt4q5lc36ri7vpcmoz4i"; const OWNER: &str = "did:plc:scallop"; const REPO: &str = "did:plc:limpet"; fn did(s: &str) -> Did { Did::new_owned(s).unwrap() } fn range(sl: u32, sc: u32, el: u32, ec: u32) -> Range { Range { start: Location { line_number: sl, column: sc, }, end: Location { line_number: el, column: ec, }, } } fn params(q: &str) -> search_code::SearchCode { search_code::SearchCode { q: q.into(), repo: None, limit: None, cursor: None, } } /// A zoekt `FileMatch` with one content chunk. `Content` is base64 on the wire, since the /// Go field is a `[]byte`. fn content_match(name: &str, line: u32, content: &str, cols: (u32, u32)) -> Value { use base64::Engine as _; json!({ "FileName": name, "Repository": "repo-7", "Version": "abc123", "Language": "Rust", "Branches": ["main"], "ChunkMatches": [{ "Content": base64::engine::general_purpose::STANDARD.encode(content), "ContentStart": {"LineNumber": line, "Column": 1}, "Ranges": [{ "Start": {"LineNumber": line, "Column": cols.0}, "End": {"LineNumber": line, "Column": cols.1}, }], }], }) } async fn mock_zoekt(server: &MockServer, files: Value) { Mock::given(method("POST")) .and(path("/api/search")) .respond_with(ResponseTemplate::new(200).set_body_json(json!({ "Result": { "Files": files, "RepoURLs": { "repo-7": format!("https://tangled.org/{REPO}/blob/{{{{.Version}}}}/{{{{.Path}}}}"), }, } }))) .mount(server) .await; } /// Slingshot and zoekt both point at `server`. When `hydratable`, the repo record is /// mounted and observed so `build_repo_view_basic` resolves. async fn app(server: &MockServer, hydratable: bool) -> AppState { let base = Url::parse(&server.uri()).unwrap(); if hydratable { Mock::given(method("GET")) .and(path("/xrpc/com.atproto.repo.getRecord")) .and(query_param("repo", OWNER)) .and(query_param("collection", "sh.tangled.repo")) .and(query_param("rkey", "r1")) .respond_with(ResponseTemplate::new(200).set_body_json(json!({ "uri": format!("at://{OWNER}/sh.tangled.repo/r1"), "cid": CID, "value": { "$type": "sh.tangled.repo", "name": "core", "knot": "oyster.cafe", "createdAt": "2026-05-01T00:00:00Z", "repoDid": REPO, }, }))) .mount(server) .await; } let state = AppState::new( Arc::new(LruRecordStore::new(CacheCapacity::from_bytes(64 * 1024))), SlingshotClient::with_default_http(base.clone()).unwrap(), Arc::new(EdgeStore::new(RuntimeHasher::default())), Arc::new(StateIndex::new(RuntimeHasher::default())), Arc::new(StateIndex::new(RuntimeHasher::default())), Arc::new(CoverageWatch::new()), Arc::new( KnotProxy::new( KnotProxyConfig::default(), KnotHttpConfig::default(), Arc::new(SystemClock::new()), RuntimeHasher::default(), ) .unwrap(), ), Arc::new( SearchIndex::new(DEFAULT_WRITER_HEAP_BYTES, Arc::new(SystemClock::new())).unwrap(), ) as Arc, Arc::new(RepoIdResolver::detached(RuntimeHasher::default())), Arc::new(crate::default_directory()), ); if hydratable { state .resolver .observe( did(OWNER), Rkey::new_owned("r1").unwrap(), Some(did(REPO)), None, &Datetime::raw_str("2026-05-01T00:00:00Z"), ) .await; } state.with_codesearch(Some(Arc::new(CodeSearch::new(&base).unwrap()))) } async fn run(state: &AppState, p: search_code::SearchCode) -> Result { search_code_inner(state, None, p) .await .map(|resp| serde_json::to_value(resp.0).unwrap()) } #[tokio::test] async fn answers_not_implemented_when_zoekt_is_unconfigured() { let server = MockServer::start().await; let state = app(&server, true).await.with_codesearch(None); let err = run(&state, params("fn main")).await.expect_err("must 501"); assert!(matches!(err, XrpcError::NotImplemented(_)), "got {err}"); } #[tokio::test] async fn returns_hydrated_file_results() { let server = MockServer::start().await; // "é" is two bytes, so the rune columns zoekt reports are not byte offsets mock_zoekt( &server, json!([content_match("src/main.rs", 12, "// héllo\n", (4, 9))]), ) .await; let state = app(&server, true).await; let body = run(&state, params("fn main")).await.expect("ok"); let result = &body["results"][0]; assert_eq!(result["path"], "src/main.rs"); assert_eq!(result["commit"], "abc123"); assert_eq!(result["language"], "Rust"); assert_eq!(result["branches"][0], "main"); // the repo view came out of the mounted record, keyed by the DID scraped from RepoURLs assert_eq!(result["repo"]["did"], REPO); assert_eq!(result["repo"]["slug"], "core"); assert_eq!(result["repo"]["owner"]["did"], OWNER); // base64 decoded, and rune columns 4..9 landed on bytes 3..9 let chunk = &result["chunks"][0]; assert_eq!(chunk["content"], "// héllo\n"); assert_eq!(chunk["lineStart"], 12); assert_eq!(chunk["highlights"][0]["start"], 3); assert_eq!(chunk["highlights"][0]["end"], 9); // one page, so no cursor assert!(body.get("cursor").is_none(), "got {body}"); } #[tokio::test] async fn scopes_by_repo_did_and_pages_with_an_offset_cursor() { let server = MockServer::start().await; // three docs for a limit of one: zoekt is asked for offset+limit+1, and the extra doc // is what tells us a further page exists mock_zoekt( &server, json!([ content_match("a.rs", 1, "one\n", (1, 4)), content_match("b.rs", 1, "two\n", (1, 4)), content_match("c.rs", 1, "three\n", (1, 4)), ]), ) .await; let state = app(&server, true).await; let body = run( &state, search_code::SearchCode { q: "fn main".into(), repo: Some(did(REPO)), limit: Some(1), cursor: Some("1".into()), }, ) .await .expect("ok"); // windowed to the second doc assert_eq!(body["results"].as_array().unwrap().len(), 1); assert_eq!(body["results"][0]["path"], "b.rs"); assert_eq!(body["cursor"], "2"); let requests = server.received_requests().await.unwrap(); let sent: Value = serde_json::from_slice( &requests .iter() .find(|r| r.url.path() == "/api/search") .expect("zoekt was called") .body, ) .unwrap(); assert_eq!(sent["Q"], format!("meta.did:{REPO} fn main")); assert!(sent["Opts"]["ChunkMatches"].as_bool().unwrap()); assert_eq!(sent["Opts"]["MaxDocDisplayCount"], 3); // offset 1 + limit 1 + probe } #[tokio::test] async fn drops_results_whose_repo_cannot_be_hydrated() { let server = MockServer::start().await; mock_zoekt(&server, json!([content_match("a.rs", 1, "one\n", (1, 4))])).await; // no repo record mounted and nothing observed, so hydration misses let state = app(&server, false).await; let body = run(&state, params("fn main")).await.expect("ok"); assert!(body["results"].as_array().unwrap().is_empty(), "got {body}"); } #[tokio::test] async fn rejects_a_blank_query_and_a_non_numeric_cursor() { let server = MockServer::start().await; let state = app(&server, true).await; let err = run(&state, params(" ")).await.expect_err("must reject"); assert!(matches!(err, XrpcError::InvalidParams(_)), "got {err}"); let mut p = params("fn"); p.cursor = Some("abc".into()); let err = run(&state, p).await.expect_err("must reject"); match err { XrpcError::InvalidParams(m) => assert!(m.contains("cursor"), "got {m}"), other => panic!("got {other}"), } } #[tokio::test] async fn maps_a_zoekt_query_rejection_to_invalid_params() { let server = MockServer::start().await; Mock::given(method("POST")) .and(path("/api/search")) .respond_with(ResponseTemplate::new(400).set_body_string("parse error")) .mount(&server) .await; let state = app(&server, true).await; let err = run(&state, params("(unbalanced")) .await .expect_err("must reject"); assert!(matches!(err, XrpcError::InvalidParams(_)), "got {err}"); } #[test] fn maps_rune_columns_to_byte_offsets() { // "é" is two bytes, so rune columns on line 2 must not be read as byte offsets let content = "// héllo\nlet x = 1;\n"; let highlights = chunk_highlights( content, 7, &[ // line 7, runes 4..9 => "héllo", bytes 3..9 (one extra byte for é) range(7, 4, 7, 9), // line 8, runes 5..6 => "x", and line 8 starts at byte 10 range(8, 5, 8, 6), ], ); let spans: Vec<&str> = highlights .iter() .map(|h| &content[h.start as usize..h.end as usize]) .collect(); assert_eq!(spans, vec!["héllo", "x"]); assert_eq!((highlights[0].start, highlights[0].end), (3, 9)); } #[test] fn drops_ranges_outside_the_chunk() { let content = "let x = 1;\n"; // before the chunk, after the chunk, and empty let highlights = chunk_highlights( content, 7, &[range(6, 1, 6, 4), range(9, 1, 9, 4), range(7, 3, 7, 3)], ); assert!(highlights.is_empty(), "got {highlights:?}"); } #[test] fn clamps_a_column_past_the_end_of_the_line() { let content = "ab\n"; // rune column 99 has no byte; clamp to the offset before the newline let highlights = chunk_highlights(content, 1, &[range(1, 1, 1, 99)]); assert_eq!((highlights[0].start, highlights[0].end), (0, 2)); } }