g1t/services/search/src/lib.rs
Pick any line to see why it is the way it is: the commit, the pull request and issue it came from, and what the agent was thinking.
| Search across all of g1t, Explore, and a command palette | 1 | //! The search service: one search across all of g1t, and Explore. |
| 2 | //! | |
| 3 | //! - **What it searches.** Repositories (name, description, topics, the | |
| 4 | //! opening of the README), code on default branches (paths and | |
| 5 | //! contents), issues and pull requests (titles and bodies), people and | |
| 6 | //! workspaces. Everything public, and private content the viewer is a | |
| 7 | //! member of. | |
| 8 | //! - **The index.** D1 with FTS5, a database that costs nothing while no | |
| 9 | //! one uses it: words for prose (`unicode61`), every run of three | |
| 10 | //! characters for code (`trigram`), each table kept in step with its | |
| 11 | //! index by triggers. See `migrations/0001_init.sql`. | |
| 12 | //! - **Visibility.** Decided when a query runs, twice: in SQL against the | |
| 13 | //! viewer's memberships now, then on the page against the repos service. | |
| 14 | //! See `visibility.rs`. | |
| 15 | //! - **Keeping it current.** From the events bus (`g1t-events-search`) and | |
| 16 | //! its own job queue (`g1t-search-jobs`), with the cost of each push | |
| 17 | //! capped. See `index.rs`. | |
| 18 | //! | |
| 19 | //! Reached through service bindings: `POST /rpc/<method>`; see | |
| 20 | //! `g1t_contracts::search`. | |
| 21 | ||
| 22 | mod index; | |
| 23 | mod query; | |
| 24 | mod read; | |
| 25 | mod rules; | |
| 26 | mod snippet; | |
| 27 | mod sql; | |
| 28 | mod store; | |
| 29 | mod visibility; | |
| 30 | ||
| 31 | use std::cell::RefCell; | |
| 32 | use std::collections::HashMap; | |
| 33 | ||
| 34 | use g1t_contracts::User; | |
| 35 | use g1t_contracts::events::Event; | |
| 36 | use g1t_kit::{args, reply, rpc_method}; | |
| 37 | use worker::{Context, D1Database, Env, Fetcher, MessageBatch, MessageExt, Request, Response, Result, event}; | |
| 38 | ||
| 39 | use index::Job; | |
| 40 | ||
| 41 | /// The queue of this service's own jobs. | |
| 42 | const JOBS_QUEUE: &str = "g1t-search-jobs"; | |
| 43 | ||
| 44 | pub struct Search { | |
| 45 | db: D1Database, | |
| 46 | env: Env, | |
| 47 | repos: Fetcher, | |
| 48 | work: Fetcher, | |
| 49 | identity: Fetcher, | |
| 50 | /// Each workspace's own principal, asked of identity once per invocation. | |
| 51 | actors: RefCell<HashMap<String, Option<User>>>, | |
| 52 | } | |
| 53 | ||
| 54 | impl Search { | |
| 55 | fn new(env: Env) -> Result<Self> { | |
| 56 | Ok(Search { | |
| 57 | db: env.d1("DB")?, | |
| 58 | repos: env.service("REPOS")?, | |
| 59 | work: env.service("WORK")?, | |
| 60 | identity: env.service("IDENTITY")?, | |
| 61 | actors: RefCell::new(HashMap::new()), | |
| 62 | env, | |
| 63 | }) | |
| 64 | } | |
| 65 | } | |
| 66 | ||
| 67 | #[event(fetch)] | |
| 68 | async fn fetch(mut request: Request, env: Env, _ctx: Context) -> Result<Response> { | |
| 69 | let Some(method) = rpc_method(&request) else { | |
| 70 | return Response::error("Not found", 404); | |
| 71 | }; | |
| 72 | let search = Search::new(env)?; | |
| 73 | let body: serde_json::Value = request.json().await?; | |
| 74 | match method.as_str() { | |
| 75 | "search" => reply(&search.search(args(body)?).await?), | |
| 76 | "suggest" => reply(&search.suggest(args(body)?).await?), | |
| 77 | "explore" => reply(&search.explore(args(body)?).await?), | |
| 78 | _ => Response::error("Unknown method", 404), | |
| 79 | } | |
| 80 | } | |
| 81 | ||
| 82 | /// Events from the bus, and this service's own jobs. Indexing is | |
| 83 | /// idempotent, so a message that fails is retried. | |
| 84 | #[event(queue)] | |
| 85 | async fn queue(batch: MessageBatch<serde_json::Value>, env: Env, _ctx: Context) -> Result<()> { | |
| 86 | let search = Search::new(env)?; | |
| 87 | let jobs = batch.queue() == JOBS_QUEUE; | |
| 88 | if !jobs && let Err(error) = search.ensure_backfill().await { | |
| 89 | worker::console_error!("search: could not start the backfill: {error}"); | |
| 90 | } | |
| 91 | for message in batch.messages()? { | |
| 92 | let body = message.body().clone(); | |
| 93 | let outcome = if jobs { | |
| 94 | match serde_json::from_value::<Job>(body) { | |
| 95 | Ok(job) => search.run_job(job).await, | |
| 96 | Err(error) => { | |
| 97 | worker::console_error!("search: a job could not be read: {error}"); | |
| 98 | Ok(()) | |
| 99 | } | |
| 100 | } | |
| 101 | } else { | |
| 102 | match serde_json::from_value::<Event>(body) { | |
| 103 | Ok(event) => search.on_event(&event).await, | |
| 104 | Err(_) => Ok(()), | |
| 105 | } | |
| 106 | }; | |
| 107 | match outcome { | |
| 108 | Ok(()) => message.ack(), | |
| 109 | Err(error) => { | |
| 110 | worker::console_error!("search: a message failed and will be retried: {error}"); | |
| 111 | message.retry(); | |
| 112 | } | |
| 113 | } | |
| 114 | } | |
| 115 | Ok(()) | |
| 116 | } |