Skip to content
158 linesCodeBlameRaw
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
22mod index;
23mod query;
24mod read;
25mod rules;
26mod snippet;
27mod sql;
28mod store;
29mod visibility;
30
31use std::cell::RefCell;
32use std::collections::HashMap;
33
34use g1t_contracts::User;
35use g1t_contracts::billing::PauseLevel;
36use g1t_contracts::events::Event;
37use g1t_kit::{args, reply, rpc_method};
38use worker::{Context, D1Database, Env, Fetcher, MessageBatch, MessageExt, Request, Response, Result, event};
39
40use index::Job;
41
42/// The queue of this service's own jobs.
43const JOBS_QUEUE: &str = "g1t-search-jobs";
44
45/// How often an isolate looks for backfill pages parked while indexing
46/// was paused.
47const RESUME_EVERY_MS: u64 = 5 * 60 * 1000;
48
49thread_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.
54fn 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
64pub 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
74impl 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)]
88async 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 };
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;
96 let body: serde_json::Value = request.json().await?;
97 let answered = match method.as_str() {
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),
102 };
103 served.finish(answered)
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)]
109async 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;
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 }
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) {
136 Ok(job) if paused && index::parked_key(&job).is_some() => search.park(&job).await,
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}