diff --git a/Cargo.lock b/Cargo.lock index c6efab97d..21ae47ba8 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -775,6 +775,21 @@ dependencies = [ "url", ] +[[package]] +name = "bobbin-codesearch" +version = "0.0.1" +dependencies = [ + "base64", + "http", + "jacquard-common", + "reqwest 0.13.1", + "serde", + "serde_json", + "thiserror 2.0.18", + "tracing", + "url", +] + [[package]] name = "bobbin-edge-index" version = "0.0.1" @@ -992,6 +1007,7 @@ version = "0.0.1" dependencies = [ "axum", "base64", + "bobbin-codesearch", "bobbin-edge-index", "bobbin-ingest", "bobbin-knot-proxy", @@ -1007,6 +1023,7 @@ dependencies = [ "jacquard-axum", "jacquard-common", "jacquard-identity", + "lexicons", "reqwest 0.13.1", "scc", "serde", diff --git a/Cargo.toml b/Cargo.toml index 9bbf0d8fc..d555b419c 100644 --- a/Cargo.toml +++ b/Cargo.toml @@ -33,6 +33,7 @@ lexicons = { path = "crates/lexicons" } trusted-proxies = { path = "crates/trusted-proxies" } bobbin-types = { path = "bobbin/crates/types" } +bobbin-codesearch = { path = "bobbin/crates/codesearch" } bobbin-edge-index = { path = "bobbin/crates/edge-index" } bobbin-ingest = { path = "bobbin/crates/ingest" } bobbin-resolver = { path = "bobbin/crates/resolver" } diff --git a/bobbin/crates/bobbin/src/config.rs b/bobbin/crates/bobbin/src/config.rs index 911de1064..b90bdfbd9 100644 --- a/bobbin/crates/bobbin/src/config.rs +++ b/bobbin/crates/bobbin/src/config.rs @@ -29,6 +29,7 @@ const KNOWN_KEYS: &[&str] = &[ "service_auth.did", "record_cache.lru_bytes", "search.heap_bytes", + "codesearch.zoekt_url", "knot.allow_private", "knot.require_https", "mirror.url", @@ -55,6 +56,7 @@ const KNOWN_ENVS: &[&str] = &[ "BOBBIN_SERVICE_DID", "BOBBIN_RECORD_LRU_BYTES", "BOBBIN_SEARCH_HEAP_BYTES", + "BOBBIN_CODESEARCH_ZOEKT_URL", "BOBBIN_KNOT_ALLOW_PRIVATE", "BOBBIN_KNOT_REQUIRE_HTTPS", "BOBBIN_MIRROR_URL", @@ -89,6 +91,9 @@ pub struct BobbinConfig { #[config(nested)] pub search: SearchConfig, + #[config(nested)] + pub codesearch: CodeSearchConfig, + #[config(nested)] pub knot: KnotConfig, @@ -262,6 +267,14 @@ pub struct SearchConfig { pub heap_bytes: u64, } +#[derive(Debug, Config)] +pub struct CodeSearchConfig { + /// Origin of a zoekt-webserver with its JSON api enabled. + /// Unset disables `org.tangled.temp.search.searchCode`. + #[config(env = "BOBBIN_CODESEARCH_ZOEKT_URL")] + pub zoekt_url: Option, +} + #[derive(Debug, Config)] pub struct KnotConfig { /// Whether to allow the knot proxy to dial private/loopback addresses. Off in diff --git a/bobbin/crates/bobbin/src/main.rs b/bobbin/crates/bobbin/src/main.rs index f2f2d22f9..4258a723e 100644 --- a/bobbin/crates/bobbin/src/main.rs +++ b/bobbin/crates/bobbin/src/main.rs @@ -22,7 +22,7 @@ use bobbin_search::{SearchIndex, SearchReader}; use bobbin_slingshot_client::SlingshotClient; use bobbin_slingshot_client::default_http_client; use bobbin_xrpc::{ - AppState, HeavyLimiter, MaxInFlight, PerRequestAnonBytes, ReservedFloor, router, + AppState, CodeSearch, HeavyLimiter, MaxInFlight, PerRequestAnonBytes, ReservedFloor, router, }; use clap::{Parser, Subcommand}; use jacquard_common::deps::fluent_uri::Uri; @@ -262,6 +262,19 @@ async fn run(cfg: BobbinConfig) -> anyhow::Result<()> { tracing::info!("some sh.tangled.git.* methods will fail, since mirror_v2.url is unset",) } } + let codesearch = cfg + .codesearch + .zoekt_url + .as_ref() + .map(|url| CodeSearch::new(url).map(Arc::new)) + .transpose() + .context("codesearch.zoekt_url")?; + match cfg.codesearch.zoekt_url.as_ref() { + Some(url) => tracing::info!(zoekt = %url, "code search enabled"), + None => tracing::info!( + "org.tangled.temp.search.searchCode will fail, since codesearch.zoekt_url is unset", + ), + } let search_heap = usize::try_from(search_heap_cap) .with_context(|| format!("search heap {search_heap_cap} exceeds usize"))?; let search = Arc::new(SearchIndex::new(search_heap, clock.clone())?); @@ -386,6 +399,7 @@ async fn run(cfg: BobbinConfig) -> anyhow::Result<()> { resolver, directory, ) + .with_codesearch(codesearch) .with_identity(identity) .with_limiter(limiter) .with_mirror(mirror) diff --git a/bobbin/crates/codesearch/Cargo.toml b/bobbin/crates/codesearch/Cargo.toml new file mode 100644 index 000000000..6b4a2022a --- /dev/null +++ b/bobbin/crates/codesearch/Cargo.toml @@ -0,0 +1,18 @@ +[package] +name = "bobbin-codesearch" +version.workspace = true +edition.workspace = true +license.workspace = true +rust-version.workspace = true + +[dependencies] +jacquard-common = { workspace = true } + +base64 = { workspace = true } +http = { workspace = true } +reqwest = { workspace = true } +serde = { workspace = true } +serde_json = { workspace = true } +thiserror = { workspace = true } +tracing = { workspace = true } +url = { workspace = true } diff --git a/bobbin/crates/codesearch/src/lib.rs b/bobbin/crates/codesearch/src/lib.rs new file mode 100644 index 000000000..39e2edb61 --- /dev/null +++ b/bobbin/crates/codesearch/src/lib.rs @@ -0,0 +1,281 @@ +//! Client for zoekt-webserver's JSON search API. + +// TODO(boltless): run our own zoekt-apiserver instead + +use std::collections::BTreeMap; +use std::time::Duration; + +use base64::Engine as _; +use base64::engine::general_purpose::STANDARD as BASE64; +use http::StatusCode; +use jacquard_common::DefaultStr; +use jacquard_common::types::did::Did; +use serde::{Deserialize, Serialize}; +use thiserror::Error; +use url::Url; + +const USER_AGENT: &str = concat!("bobbin/", env!("CARGO_PKG_VERSION")); +const REQUEST_TIMEOUT: Duration = Duration::from_secs(15); +const CONNECT_TIMEOUT: Duration = Duration::from_secs(5); +const SEARCH_PATH: &str = "api/search"; +const MAX_WALL_TIME_NANOS: u64 = 10_000_000_000; +const NUM_CONTEXT_LINES: u32 = 2; +const MAX_ERROR_BODY: usize = 4096; + +#[derive(Debug, Error)] +pub enum CodeSearchError { + #[error("invalid zoekt url scheme: {0}")] + BadScheme(String), + #[error("http client build: {0}")] + Build(String), + #[error("network: {0}")] + Network(String), + #[error("zoekt rejected the query: {0}")] + BadQuery(String), + #[error("zoekt returned status {status}: {body}")] + Upstream { status: StatusCode, body: String }, + #[error("decode zoekt reply: {0}")] + Decode(String), +} + +pub struct CodeSearch { + http: reqwest::Client, + search_url: Url, +} + +impl std::fmt::Debug for CodeSearch { + fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result { + f.debug_struct("CodeSearch") + .field("search_url", &self.search_url) + .finish_non_exhaustive() + } +} + +/// One window of zoekt file matches, plus whether a following page exists. +#[derive(Debug, Default)] +pub struct ZoektPage { + pub files: Vec, + /// Repository name to `FileURLTemplate`. Resolve this to repo DID. + pub repo_urls: BTreeMap, + pub has_more: bool, +} + +impl CodeSearch { + pub fn new(base: &Url) -> Result { + match base.scheme() { + "http" | "https" => {} + other => return Err(CodeSearchError::BadScheme(other.to_owned())), + } + let search_url = base + .join(SEARCH_PATH) + .map_err(|e| CodeSearchError::BadScheme(e.to_string()))?; + let http = reqwest::Client::builder() + .user_agent(USER_AGENT) + .timeout(REQUEST_TIMEOUT) + .connect_timeout(CONNECT_TIMEOUT) + .build() + .map_err(|e| CodeSearchError::Build(e.to_string()))?; + Ok(Self { http, search_url }) + } + + pub async fn search( + &self, + query: &str, + offset: usize, + limit: usize, + ) -> Result { + let args = SearchArgs { + q: query, + opts: SearchOpts { + chunk_matches: true, + max_wall_time: MAX_WALL_TIME_NANOS, + num_context_lines: NUM_CONTEXT_LINES, + max_doc_display_count: offset.saturating_add(limit).saturating_add(1), + }, + }; + + let resp = self + .http + .post(self.search_url.clone()) + .json(&args) + .send() + .await + .map_err(|e| CodeSearchError::Network(e.to_string()))?; + + let status = resp.status(); + if !status.is_success() { + let mut body = resp.text().await.unwrap_or_default(); + body.truncate(MAX_ERROR_BODY); + let body = body.trim().to_owned(); + // zoekt answers 400 when its own query parser refuses the string + return Err(if status == StatusCode::BAD_REQUEST { + CodeSearchError::BadQuery(body) + } else { + CodeSearchError::Upstream { status, body } + }); + } + + let bytes = resp + .bytes() + .await + .map_err(|e| CodeSearchError::Network(e.to_string()))?; + let reply: SearchReply = + serde_json::from_slice(&bytes).map_err(|e| CodeSearchError::Decode(e.to_string()))?; + + let Some(result) = reply.result else { + return Ok(ZoektPage::default()); + }; + let all = result.files.unwrap_or_default(); + let end = offset.saturating_add(limit); + Ok(ZoektPage { + has_more: all.len() > end, + files: all.into_iter().skip(offset).take(limit).collect(), + repo_urls: result.repo_urls.unwrap_or_default(), + }) + } +} + +// HACK(boltless): obviously this is a hack. We should run our own zoekt api server that uses DID +// as an identifier. +/// pulls the repo DID from zoekt `FileURLTemplate` shaped +/// `{appviewURL}/{repoDID}/blob/{commit}/{path}`. +pub fn extract_did(template: &str) -> Option> { + let url = Url::parse(template).ok()?; + let seg = url.path_segments()?.next()?; + Did::new_owned(seg).ok() +} + +#[derive(Serialize)] +#[serde(rename_all = "PascalCase")] +struct SearchArgs<'a> { + q: &'a str, + opts: SearchOpts, +} + +#[derive(Serialize)] +#[serde(rename_all = "PascalCase")] +struct SearchOpts { + chunk_matches: bool, + max_wall_time: u64, + num_context_lines: u32, + max_doc_display_count: usize, +} + +#[derive(Deserialize)] +#[serde(rename_all = "PascalCase")] +struct SearchReply { + result: Option, +} + +#[derive(Deserialize)] +#[serde(rename_all = "PascalCase")] +struct SearchResult { + files: Option>, + #[serde(rename = "RepoURLs")] + repo_urls: Option>, +} + +#[derive(Debug, Deserialize)] +#[serde(rename_all = "PascalCase")] +pub struct FileMatch { + pub file_name: String, + pub repository: String, + #[serde(default)] + pub version: String, + #[serde(default)] + pub language: String, + #[serde(default)] + pub branches: Vec, + #[serde(default)] + pub chunk_matches: Vec, +} + +#[derive(Debug, Deserialize)] +#[serde(rename_all = "PascalCase")] +pub struct ChunkMatch { + #[serde(default, deserialize_with = "base64_lossy_string")] + pub content: String, + #[serde(default)] + pub ranges: Vec, + /// True when the match is on the file's name rather than its content. + #[serde(default)] + pub file_name: bool, + #[serde(default)] + pub content_start: Location, +} + +#[derive(Debug, Deserialize)] +#[serde(rename_all = "PascalCase")] +pub struct Range { + pub start: Location, + pub end: Location, +} + +#[derive(Debug, Default, Deserialize)] +#[serde(rename_all = "PascalCase")] +pub struct Location { + pub line_number: u32, + pub column: u32, +} + +fn base64_lossy_string<'de, D: serde::Deserializer<'de>>(d: D) -> Result { + let raw = String::deserialize(d)?; + let bytes = BASE64.decode(&raw).map_err(serde::de::Error::custom)?; + // lossy, not strict: one stray byte in one indexed file must not fail the whole page + Ok(String::from_utf8_lossy(&bytes).into_owned()) +} + +#[cfg(test)] +mod tests { + use super::*; + + #[test] + fn decodes_a_zoekt_reply() { + let reply: SearchReply = serde_json::from_str( + r#"{"Result":{"Files":[{ + "FileName":"main.rs","Repository":"repo-7","Version":"abc123", + "Language":"Rust","Branches":["main"], + "ChunkMatches":[ + {"Content":"Zm4gbWFpbigp","ContentStart":{"LineNumber":12,"Column":1}, + "Ranges":[{"Start":{"LineNumber":12,"Column":4},"End":{"LineNumber":12,"Column":8}}]}, + {"FileName":true,"Ranges":[{"Start":{"LineNumber":1,"Column":1},"End":{"LineNumber":1,"Column":5}}]} + ]}], + "RepoURLs":{"repo-7":"https://tangled.org/did:plc:abc/blob/{{.Version}}/{{.Path}}"}}}"#, + ) + .expect("canned reply must decode"); + + let result = reply.result.expect("Result present"); + let files = result.files.expect("Files present"); + let chunks = &files[0].chunk_matches; + // base64 "Zm4gbWFpbigp" decodes to the source line + assert_eq!(chunks[0].content, "fn main()"); + assert_eq!(chunks[0].content_start.line_number, 12); + assert!(!chunks[0].file_name); + // a filename match carries ranges but no content + assert!(chunks[1].file_name); + assert_eq!(chunks[1].content, ""); + assert_eq!(result.repo_urls.expect("RepoURLs present").len(), 1); + } + + #[test] + fn tolerates_null_files_and_repo_urls() { + // no omitempty on either field, so an empty result is `null`, not absent + let reply: SearchReply = + serde_json::from_str(r#"{"Result":{"Files":null,"RepoURLs":null}}"#).unwrap(); + let result = reply.result.unwrap(); + assert!(result.files.unwrap_or_default().is_empty()); + assert!(result.repo_urls.unwrap_or_default().is_empty()); + } + + #[test] + fn extracts_did_from_url_template() { + assert_eq!( + extract_did("https://tangled.org/did:plc:abc123/blob/{{.Version}}/{{.Path}}") + .expect("did present") + .as_str(), + "did:plc:abc123", + ); + assert!(extract_did("").is_none()); + assert!(extract_did("https://tangled.org/not-a-did/blob/x/y").is_none()); + } +} diff --git a/bobbin/crates/xrpc/Cargo.toml b/bobbin/crates/xrpc/Cargo.toml index 38d004a89..fafc5d695 100644 --- a/bobbin/crates/xrpc/Cargo.toml +++ b/bobbin/crates/xrpc/Cargo.toml @@ -7,6 +7,8 @@ rust-version.workspace = true [dependencies] bobbin-types = { workspace = true } +bobbin-codesearch = { workspace = true } +lexicons = { workspace = true, features = ["org_tangled"] } bobbin-edge-index = { workspace = true } bobbin-record-lru = { workspace = true } bobbin-resolver = { workspace = true } diff --git a/bobbin/crates/xrpc/src/codesearch.rs b/bobbin/crates/xrpc/src/codesearch.rs new file mode 100644 index 000000000..b068aec49 --- /dev/null +++ b/bobbin/crates/xrpc/src/codesearch.rs @@ -0,0 +1,526 @@ +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::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, + ) + .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)); + } +} diff --git a/bobbin/crates/xrpc/src/lib.rs b/bobbin/crates/xrpc/src/lib.rs index e38d277f3..d011abffc 100644 --- a/bobbin/crates/xrpc/src/lib.rs +++ b/bobbin/crates/xrpc/src/lib.rs @@ -22,6 +22,7 @@ use axum::{ response::{IntoResponse, Json, Response}, routing::get, }; +pub use bobbin_codesearch::{CodeSearch, CodeSearchError}; use bobbin_edge_index::{ Coverage, CoverageWatch, CursorParseError, EdgeItem, EdgePage, EdgeStore, IssueStateKind, PageCursor, PageLimit, PageOffset, PageStart, PageToken, ParsedCid, PullStatusKind, Rejection, @@ -110,6 +111,7 @@ use tracing::{Level, Span}; mod actor; mod backpressure; mod client_address; +mod codesearch; mod enrich; mod feed; mod filter; @@ -150,6 +152,8 @@ pub struct AppState { pub mirror: Option>, pub mirror_v2: Option>, pub search: Arc, + /// Zoekt-backed code search. `None` when unconfigured, and the endpoint answers 501. + pub codesearch: Option>, pub resolver: Arc, pub identity: Arc, pub directory: Arc, @@ -192,6 +196,7 @@ impl AppState { mirror: None, mirror_v2: None, search, + codesearch: None, resolver, identity, directory: directory.clone(), @@ -228,6 +233,11 @@ impl AppState { self } + pub fn with_codesearch(mut self, codesearch: Option>) -> Self { + self.codesearch = codesearch; + self + } + pub fn with_identity(mut self, identity: Arc) -> Self { self.identity = identity; self @@ -530,6 +540,10 @@ pub fn router(state: AppState) -> Router { .route("/xrpc/sh.tangled.string.listStrings", get(list_strings)) .route("/xrpc/sh.tangled.string.countStrings", get(count_strings)) .route("/xrpc/sh.tangled.search.query", get(search_query)) + .route( + "/xrpc/org.tangled.temp.search.searchCode", + get(codesearch::search_code), + ) .route( "/xrpc/sh.tangled.query.enrichResponse", axum::routing::post(enrich::enrich), @@ -889,6 +903,8 @@ pub enum XrpcError { Overloaded, #[error("record has not been ingested yet")] NotSettled, + #[error("not implemented: {0}")] + NotImplemented(&'static str), } impl XrpcError { @@ -915,6 +931,7 @@ impl IntoResponse for XrpcError { Self::Internal(_) => (StatusCode::INTERNAL_SERVER_ERROR, "InternalError"), Self::Overloaded => (StatusCode::SERVICE_UNAVAILABLE, "Overloaded"), Self::NotSettled => (StatusCode::GATEWAY_TIMEOUT, "NotSettled"), + Self::NotImplemented(_) => (StatusCode::NOT_IMPLEMENTED, "MethodNotImplemented"), }; let body = ErrorBody { error, @@ -1778,7 +1795,8 @@ fn drop_unhydratable( err @ (XrpcError::Internal(_) | XrpcError::Overloaded | XrpcError::NotSettled - | XrpcError::AuthRequired(_)), + | XrpcError::AuthRequired(_) + | XrpcError::NotImplemented(_)), ) => Err(err), } } diff --git a/bobbin/crates/xrpc/tests/code_search.rs b/bobbin/crates/xrpc/tests/code_search.rs new file mode 100644 index 000000000..389ea1eff --- /dev/null +++ b/bobbin/crates/xrpc/tests/code_search.rs @@ -0,0 +1,87 @@ +use std::sync::Arc; + +use axum::body::Body; +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 bobbin_xrpc::{AppState, CodeSearch, router}; +use http::{Request, StatusCode}; +use tower::ServiceExt; +use url::Url; + +const NSID: &str = "org.tangled.temp.search.searchCode"; + +fn app() -> axum::Router { + let base = Url::parse("http://127.0.0.1:1").unwrap(); + 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(bobbin_xrpc::default_directory()), + ) + .with_codesearch(Some(Arc::new(CodeSearch::new(&base).unwrap()))); + router(state) +} + +fn get(auth: Option<&str>) -> Request { + let mut req = Request::builder().uri(format!("/xrpc/{NSID}?q=fn+main")); + if let Some(token) = auth { + req = req.header("authorization", token); + } + req.body(Body::empty()).unwrap() +} + +/// Code search is gated on login, matching the appview ui it backs, so an anonymous request +/// must never reach zoekt. +#[tokio::test] +async fn an_anonymous_request_is_rejected() { + let resp = app().oneshot(get(None)).await.unwrap(); + assert_eq!(resp.status(), StatusCode::UNAUTHORIZED); +} + +#[tokio::test] +async fn a_malformed_authorization_header_is_rejected() { + for header in ["", "Bearer", "Bearer not-a-jwt", "Basic abc"] { + let resp = app().oneshot(get(Some(header))).await.unwrap(); + assert_eq!( + resp.status(), + StatusCode::UNAUTHORIZED, + "header {header:?} must not pass" + ); + } +} + +/// The lexicon declares a query, so the route must not accept a POST. +#[tokio::test] +async fn the_nsid_is_a_query_not_a_procedure() { + let resp = app() + .oneshot( + Request::builder() + .method("POST") + .uri(format!("/xrpc/{NSID}?q=fn")) + .body(Body::empty()) + .unwrap(), + ) + .await + .unwrap(); + assert_eq!(resp.status(), StatusCode::METHOD_NOT_ALLOWED); +} diff --git a/bobbin/example.toml b/bobbin/example.toml index e337a4223..751a23af6 100644 --- a/bobbin/example.toml +++ b/bobbin/example.toml @@ -126,6 +126,11 @@ # Default value: "http://127.0.0.1:13011" #url = "http://127.0.0.1:13011" +[service_auth] +# Can also be specified via environment variable `BOBBIN_SERVICE_DID`. +# Required! This value must be specified. +#did = + [record_cache] # Bound of bytes on the in-process record LRU. Records evict on a weighted # LRU policy keyed on URI plus payload length. @@ -145,6 +150,13 @@ # Default value: 50000000 #heap_bytes = 50000000 +[codesearch] +# Origin of a zoekt-webserver with its JSON api enabled. +# Unset disables `org.tangled.temp.search.searchCode`. +# +# Can also be specified via environment variable `BOBBIN_CODESEARCH_ZOEKT_URL`. +#zoekt_url = + [knot] # Whether to allow the knot proxy to dial private/loopback addresses. Off in # production - on for local testing against a knotserver on localhost. @@ -167,6 +179,8 @@ #url = [mirror_v2] +# Origin of a v2 mirror (gitmirror) XRPC server. +# # Can also be specified via environment variable `BOBBIN_MIRROR_V2_URL`. #url = diff --git a/bobbin/worker/src/index.ts b/bobbin/worker/src/index.ts index 3dcea4ac7..5dfd428a5 100644 --- a/bobbin/worker/src/index.ts +++ b/bobbin/worker/src/index.ts @@ -8,6 +8,7 @@ export interface Env { BOBBIN_MIRROR_V2_URL: string; BOBBIN_SERVICE_DID: string; BOBBIN_LOG: string; + BOBBIN_CODESEARCH_ZOEKT_URL?: string; } export class BobbinContainer extends Container { @@ -27,6 +28,9 @@ export class BobbinContainer extends Container { BOBBIN_SERVICE_DID: env.BOBBIN_SERVICE_DID, BOBBIN_LOG: env.BOBBIN_LOG, BOBBIN_LOG_FORMAT: "json", + ...(env.BOBBIN_CODESEARCH_ZOEKT_URL + ? { BOBBIN_CODESEARCH_ZOEKT_URL: env.BOBBIN_CODESEARCH_ZOEKT_URL } + : {}), }, }); } diff --git a/docker-compose.yml b/docker-compose.yml index 4dfa8058f..46f913e6c 100644 --- a/docker-compose.yml +++ b/docker-compose.yml @@ -535,6 +535,7 @@ services: BOBBIN_KNOT_ALLOW_PRIVATE: "true" BOBBIN_KNOT_REQUIRE_HTTPS: "false" BOBBIN_MIRROR_V2_URL: http://knotmirror:7000 + BOBBIN_CODESEARCH_ZOEKT_URL: https://zoekt.tngl.boltless.dev BOBBIN_LOG: info volumes: - .:/src:cached