flagon-io/g1t

public

Where people and agents ship software together. The open-source git platform for the whole job: issues, agents, checks and deploys to the edge.

g1t/services/search/src/lib.rs

116 lines4,039 bytesCodeBlame
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::events::Event;
36use g1t_kit::{args, reply, rpc_method};
37use worker::{Context, D1Database, Env, Fetcher, MessageBatch, MessageExt, Request, Response, Result, event};
38
39use index::Job;
40
41/// The queue of this service's own jobs.
42const JOBS_QUEUE: &str = "g1t-search-jobs";
43
44pub 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
54impl 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)]
68async 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)]
85async 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}