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; | |
| Merge platform pause and the hourly usage watcher: staff can pause compute, schedules, indexing or renders for everyone, the watcher emails on a breach and is never blind quietly, and the models proxy holds each run to its cap (billing 0051, integrations 0006) | 35 | use g1t_contracts::billing::PauseLevel; |
| Search across all of g1t, Explore, and a command palette | 36 | use g1t_contracts::events::Event; |
| 37 | use g1t_kit::{args, reply, rpc_method}; | |
| 38 | use worker::{Context, D1Database, Env, Fetcher, MessageBatch, MessageExt, Request, Response, Result, event}; | |
| 39 | ||
| 40 | use index::Job; | |
| 41 | ||
| 42 | /// The queue of this service's own jobs. | |
| 43 | const JOBS_QUEUE: &str = "g1t-search-jobs"; | |
| 44 | ||
| Merge platform pause and the hourly usage watcher: staff can pause compute, schedules, indexing or renders for everyone, the watcher emails on a breach and is never blind quietly, and the models proxy holds each run to its cap (billing 0051, integrations 0006) | 45 | /// How often an isolate looks for backfill pages parked while indexing |
| 46 | /// was paused. | |
| 47 | const RESUME_EVERY_MS: u64 = 5 * 60 * 1000; | |
| 48 | ||
| 49 | thread_local! { | |
| 50 | static RESUME_CHECKED: std::cell::Cell<u64> = const { std::cell::Cell::new(0) }; | |
| 51 | } | |
| 52 | ||
| 53 | /// Whether to look for parked pages now; marks it looked for. | |
| 54 | fn resume_due(now: u64) -> bool { | |
| 55 | RESUME_CHECKED.with(|last| { | |
| 56 | if now.saturating_sub(last.get()) < RESUME_EVERY_MS { | |
| 57 | return false; | |
| 58 | } | |
| 59 | last.set(now); | |
| 60 | true | |
| 61 | }) | |
| 62 | } | |
| 63 | ||
| Search across all of g1t, Explore, and a command palette | 64 | pub struct Search { |
| 65 | db: D1Database, | |
| 66 | env: Env, | |
| 67 | repos: Fetcher, | |
| 68 | work: Fetcher, | |
| 69 | identity: Fetcher, | |
| 70 | /// Each workspace's own principal, asked of identity once per invocation. | |
| 71 | actors: RefCell<HashMap<String, Option<User>>>, | |
| 72 | } | |
| 73 | ||
| 74 | impl Search { | |
| 75 | fn new(env: Env) -> Result<Self> { | |
| 76 | Ok(Search { | |
| 77 | db: env.d1("DB")?, | |
| 78 | repos: env.service("REPOS")?, | |
| 79 | work: env.service("WORK")?, | |
| 80 | identity: env.service("IDENTITY")?, | |
| 81 | actors: RefCell::new(HashMap::new()), | |
| 82 | env, | |
| 83 | }) | |
| 84 | } | |
| 85 | } | |
| 86 | ||
| 87 | #[event(fetch)] | |
| 88 | async fn fetch(mut request: Request, env: Env, _ctx: Context) -> Result<Response> { | |
| 89 | let Some(method) = rpc_method(&request) else { | |
| 90 | return Response::error("Not found", 404); | |
| 91 | }; | |
| Fast pages, required checks on the branch, self-hosted runners, honest incidents | 92 | // A replica near the caller when it asks for one (crates/kit/src/d1.rs). |
| 93 | let (db, served) = g1t_kit::d1::open(&env, "DB", &request)?; | |
| 94 | let mut search = Search::new(env)?; | |
| 95 | search.db = db; | |
| Search across all of g1t, Explore, and a command palette | 96 | let body: serde_json::Value = request.json().await?; |
| Fast pages, required checks on the branch, self-hosted runners, honest incidents | 97 | let answered = match method.as_str() { |
| Search across all of g1t, Explore, and a command palette | 98 | "search" => reply(&search.search(args(body)?).await?), |
| 99 | "suggest" => reply(&search.suggest(args(body)?).await?), | |
| 100 | "explore" => reply(&search.explore(args(body)?).await?), | |
| 101 | _ => Response::error("Unknown method", 404), | |
| Fast pages, required checks on the branch, self-hosted runners, honest incidents | 102 | }; |
| 103 | served.finish(answered) | |
| Search across all of g1t, Explore, and a command palette | 104 | } |
| 105 | ||
| 106 | /// Events from the bus, and this service's own jobs. Indexing is | |
| 107 | /// idempotent, so a message that fails is retried. | |
| 108 | #[event(queue)] | |
| 109 | async fn queue(batch: MessageBatch<serde_json::Value>, env: Env, _ctx: Context) -> Result<()> { | |
| 110 | let search = Search::new(env)?; | |
| 111 | let jobs = batch.queue() == JOBS_QUEUE; | |
| Merge platform pause and the hourly usage watcher: staff can pause compute, schedules, indexing or renders for everyone, the watcher emails on a breach and is never blind quietly, and the models proxy holds each run to its cap (billing 0051, integrations 0006) | 112 | // Indexing paused across g1t (billing's `platform_pause`, kept 30 |
| 113 | // seconds in the isolate): backfills wait, and events still keep the | |
| 114 | // index current. Without a billing binding nothing is ever paused. | |
| 115 | let paused = match search.env.service("BILLING") { | |
| 116 | Ok(billing) => g1t_kit::pause::paused(&billing, PauseLevel::Indexing).await, | |
| 117 | Err(_) => false, | |
| 118 | }; | |
| 119 | if !jobs && !paused { | |
| 120 | if let Err(error) = search.ensure_backfill().await { | |
| 121 | worker::console_error!("search: could not start the backfill: {error}"); | |
| 122 | } | |
| 123 | // Pages parked while paused, looked for at most every few minutes. | |
| 124 | if resume_due(g1t_kit::now_ms()) { | |
| 125 | match search.resume_parked().await { | |
| 126 | Ok(0) => {} | |
| 127 | Ok(found) => worker::console_log!("search: resumed {found} parked backfill pages"), | |
| 128 | Err(error) => worker::console_error!("search: could not resume parked backfill pages: {error}"), | |
| 129 | } | |
| 130 | } | |
| Search across all of g1t, Explore, and a command palette | 131 | } |
| 132 | for message in batch.messages()? { | |
| 133 | let body = message.body().clone(); | |
| 134 | let outcome = if jobs { | |
| 135 | match serde_json::from_value::<Job>(body) { | |
| Merge platform pause and the hourly usage watcher: staff can pause compute, schedules, indexing or renders for everyone, the watcher emails on a breach and is never blind quietly, and the models proxy holds each run to its cap (billing 0051, integrations 0006) | 136 | Ok(job) if paused && index::parked_key(&job).is_some() => search.park(&job).await, |
| Search across all of g1t, Explore, and a command palette | 137 | Ok(job) => search.run_job(job).await, |
| 138 | Err(error) => { | |
| 139 | worker::console_error!("search: a job could not be read: {error}"); | |
| 140 | Ok(()) | |
| 141 | } | |
| 142 | } | |
| 143 | } else { | |
| 144 | match serde_json::from_value::<Event>(body) { | |
| 145 | Ok(event) => search.on_event(&event).await, | |
| 146 | Err(_) => Ok(()), | |
| 147 | } | |
| 148 | }; | |
| 149 | match outcome { | |
| 150 | Ok(()) => message.ack(), | |
| 151 | Err(error) => { | |
| 152 | worker::console_error!("search: a message failed and will be retried: {error}"); | |
| 153 | message.retry(); | |
| 154 | } | |
| 155 | } | |
| 156 | } | |
| 157 | Ok(()) | |
| 158 | } |
This file's history is long; its oldest lines are credited to the oldest commit read.