g1t/services/search/src/index.rs

941 lines38,190 bytesCodeBlame
1//! Keeping the index current: from events as things change, from queued
2//! jobs for anything bigger than one event should do, and from a backfill
3//! of everything that was there before.
4//!
5//! The cost of a push is capped. A push lists only the files it changed
6//! (up to [`MAX_PUSH_FILES`]; more than that and the repository is compared
7//! whole instead), queues them, and each job reads at most
8//! [`DRAIN_FILES`] files and [`DRAIN_BYTES`] of text before handing the
9//! rest to the next. A repository is indexed up to [`MAX_REPO_FILES`]
10//! files; past that it is marked partial.
11
12use std::collections::{HashMap, HashSet};
13
14use g1t_contracts::events::Event;
15use g1t_contracts::identity::{DirectoryArgs, DirectoryEntry, DirectoryPage, Profile, SlugArgs, UsernameArgs, Workspace};
16use g1t_contracts::repos::{
17 AllIdsArgs, BlobText, ChangedFilesArgs, FileList, IdPage, ListFilesArgs, PathByIdArgs, ReadBlobsArgs, RepoPath, TreeArgs,
18 TreeView,
19};
20use g1t_contracts::time::rfc3339;
21use g1t_contracts::work::{
22 Issue, IssueDetail, ListIssuesArgs, ListPullsArgs, Pull, PullDetail, State, ViewArgs,
23};
24use g1t_contracts::{Membership, Outcome, PrincipalKind, User};
25use g1t_kit::now_ms;
26use serde::{Deserialize, Serialize};
27use worker::Result;
28
29use crate::Search;
30use crate::rules::{self, MAX_FILE_BYTES, SKIP_DIRS};
31use crate::sql::{Param, Sql};
32use crate::store;
33
34/// Files a push may change and still be indexed from its own list; more,
35/// and the repository is compared whole.
36pub const MAX_PUSH_FILES: u32 = 300;
37/// Files indexed in one repository, at most.
38pub const MAX_REPO_FILES: u32 = 5000;
39/// Files one job reads, at most.
40pub const DRAIN_FILES: usize = 60;
41/// Text one job reads, at most.
42pub const DRAIN_BYTES: u64 = 3 * 1024 * 1024;
43/// Blobs asked of the repos service at once.
44const READ_GROUP: usize = 20;
45/// Pieces of a file written in one statement.
46const PIECES_PER_STATEMENT: usize = 8;
47/// The most of an issue's or pull request's body indexed.
48const BODY_CHARS: usize = 32 * 1024;
49/// The most of a README indexed.
50const README_CHARS: usize = 4000;
51/// Repositories listed per backfill job.
52const BACKFILL_PAGE: u32 = 100;
53/// Accounts or workspaces per directory job.
54const DIRECTORY_PAGE: u32 = 200;
55
56/// Work queued for later, one job per message.
57#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
58#[serde(tag = "type", rename_all = "snake_case")]
59pub enum Job {
60 /// Lists every repository a page at a time, queueing one `repo` job each.
61 Backfill {
62 #[serde(default)]
63 after: Option<String>,
64 },
65 /// Compares a repository's default branch with the index, whole, and
66 /// indexes its issues and pull requests.
67 Repo { repo_id: String },
68 /// Indexes what a push to the default branch changed.
69 Push {
70 repo_id: String,
71 #[serde(default)]
72 before: Option<String>,
73 after: String,
74 },
75 /// Reads the next files waiting for a repository.
76 Drain { repo_id: String },
77 /// Indexes every account or workspace, a page at a time.
78 Directory {
79 kind: String,
80 #[serde(default)]
81 after: Option<String>,
82 },
83}
84
85/// What this service reads of the events it handles.
86#[derive(Deserialize)]
87#[serde(rename_all = "camelCase")]
88struct RepoEvent {
89 repo_id: String,
90}
91
92#[derive(Deserialize)]
93#[serde(rename_all = "camelCase")]
94struct PushEvent {
95 repo_id: String,
96 #[serde(default)]
97 before: Option<String>,
98 after: String,
99 #[serde(default)]
100 default_branch: bool,
101}
102
103#[derive(Deserialize)]
104#[serde(rename_all = "camelCase")]
105struct NumberEvent {
106 repo_id: String,
107 number: u32,
108}
109
110#[derive(Deserialize)]
111struct UserEvent {
112 username: String,
113}
114
115#[derive(Deserialize)]
116struct WorkspaceEvent {
117 slug: String,
118}
119
120#[derive(Deserialize)]
121struct Stored {
122 path: String,
123 blob: String,
124}
125
126#[derive(Deserialize)]
127struct Waiting {
128 path: String,
129 blob: Option<String>,
130}
131
132#[derive(Deserialize)]
133struct RepoState {
134 state: String,
135}
136
137fn now() -> String {
138 rfc3339(now_ms())
139}
140
141fn p(text: &str) -> Param {
142 Param::Text(text.to_owned())
143}
144
145/// A commit id that means "nothing": git's forty zeros.
146fn is_null_commit(hash: &str) -> bool {
147 hash.is_empty() || hash.bytes().all(|b| b == b'0')
148}
149
150fn labels_of(labels: &[String]) -> String {
151 let mut out = String::from("|");
152 for label in labels {
153 out.push_str(&label.to_lowercase());
154 out.push('|');
155 }
156 out
157}
158
159fn clip(text: &str, max: usize) -> String {
160 text.chars().take(max).collect()
161}
162
163impl Search {
164 /// The workspace itself, through no one: how this service reads a
165 /// repository, private or not, to index it.
166 async fn actor(&self, slug: &str) -> Result<Option<User>> {
167 let slug = slug.to_lowercase();
168 if let Some(found) = self.actors.borrow().get(&slug) {
169 return Ok(found.clone());
170 }
171 let workspace: Option<Workspace> = g1t_kit::call(&self.identity, "get_workspace", &SlugArgs { slug: slug.clone() }).await?;
172 let actor = workspace.map(|workspace| User {
173 id: workspace.id,
174 username: workspace.slug.clone(),
175 kind: PrincipalKind::Workspace,
176 verified: true,
177 workspaces: vec![Membership::member(workspace.slug)],
178 ..User::default()
179 });
180 self.actors.borrow_mut().insert(slug, actor.clone());
181 Ok(actor)
182 }
183
184 async fn path_of(&self, repo_id: &str) -> Result<Option<RepoPath>> {
185 g1t_kit::call(&self.repos, "path_by_id", &PathByIdArgs { id: repo_id.to_owned() }).await
186 }
187
188 pub async fn enqueue(&self, jobs: Vec<Job>) -> Result<()> {
189 let queue = self.env.queue("JOBS")?;
190 for batch in jobs.chunks(100) {
191 queue.send_batch(batch.to_vec()).await?;
192 }
193 Ok(())
194 }
195
196 // ---- Repositories --------------------------------------------------------
197
198 /// Reads a repository's details and README again and stores them.
199 /// `None`, and the repository forgotten, when it is gone or a fork.
200 pub async fn refresh_repo(&self, repo_id: &str) -> Result<Option<TreeView>> {
201 let Some(path) = self.path_of(repo_id).await? else {
202 self.purge(repo_id).await?;
203 return Ok(None);
204 };
205 let Some(actor) = self.actor(&path.namespace).await? else {
206 return Ok(None);
207 };
208 let tree: Outcome<TreeView> = g1t_kit::call(
209 &self.repos,
210 "tree",
211 &TreeArgs { path, viewer: Some(actor), git_ref: None, tree_path: String::new() },
212 )
213 .await?;
214 let Outcome::Ok(tree) = tree else {
215 return Ok(None);
216 };
217 let repo = &tree.repo;
218 let readme = tree.readme.as_ref().and_then(|readme| readme.text.as_deref()).map(|text| clip(text, README_CHARS));
219 let topics = repo.topics.join(" ");
220 let mut statements = vec![store::prepare(
221 &self.db,
222 "INSERT INTO repos (repo_id, namespace, name, description, topics, readme, private, default_branch, created_at, updated_at)
223 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
224 ON CONFLICT (repo_id) DO UPDATE SET namespace = excluded.namespace, name = excluded.name,
225 description = excluded.description, topics = excluded.topics, readme = excluded.readme,
226 private = excluded.private, default_branch = excluded.default_branch, updated_at = excluded.updated_at",
227 vec![
228 p(&repo.id),
229 p(&repo.namespace.to_lowercase()),
230 p(&repo.name.to_lowercase()),
231 Param::opt(repo.description.as_deref()),
232 p(&topics),
233 Param::opt(readme.as_deref()),
234 Param::Int(i64::from(repo.is_private)),
235 p(&repo.default_branch),
236 p(&repo.created_at),
237 p(&now()),
238 ],
239 )?];
240 statements.push(store::prepare(&self.db, "DELETE FROM repo_topics WHERE repo_id = ?", vec![p(&repo.id)])?);
241 for topic in &repo.topics {
242 statements.push(store::prepare(
243 &self.db,
244 "INSERT OR IGNORE INTO repo_topics (topic, repo_id) VALUES (?, ?)",
245 vec![p(topic), p(&repo.id)],
246 )?);
247 }
248 store::run_all(&self.db, statements).await?;
249 Ok(Some(tree))
250 }
251
252 /// Forgets a repository and everything indexed from it.
253 pub async fn purge(&self, repo_id: &str) -> Result<()> {
254 let id = || vec![p(repo_id)];
255 store::run_all(
256 &self.db,
257 vec![
258 store::prepare(&self.db, "DELETE FROM chunks WHERE fid IN (SELECT fid FROM files WHERE repo_id = ?)", id())?,
259 store::prepare(&self.db, "DELETE FROM files WHERE repo_id = ?", id())?,
260 store::prepare(&self.db, "DELETE FROM pending WHERE repo_id = ?", id())?,
261 store::prepare(&self.db, "DELETE FROM items WHERE repo_id = ?", id())?,
262 store::prepare(&self.db, "DELETE FROM repo_topics WHERE repo_id = ?", id())?,
263 store::prepare(&self.db, "DELETE FROM repos WHERE repo_id = ?", id())?,
264 ],
265 )
266 .await
267 }
268
269 /// Queues files to read (or, with no blob, to remove).
270 async fn queue_files(&self, repo_id: &str, files: &[(String, Option<String>)]) -> Result<()> {
271 let at = now();
272 let mut statements = Vec::new();
273 for group in files.chunks(20) {
274 let rows = vec!["(?, ?, ?, ?)"; group.len()].join(", ");
275 let mut params = Vec::with_capacity(group.len() * 4);
276 for (path, blob) in group {
277 params.extend([p(repo_id), p(path), Param::opt(blob.as_deref()), p(&at)]);
278 }
279 statements.push(store::prepare(
280 &self.db,
281 &format!(
282 "INSERT INTO pending (repo_id, path, blob, queued_at) VALUES {rows}
283 ON CONFLICT (repo_id, path) DO UPDATE SET blob = excluded.blob, queued_at = excluded.queued_at"
284 ),
285 params,
286 )?);
287 }
288 store::run_all(&self.db, statements).await
289 }
290
291 fn skip_dirs() -> Vec<String> {
292 SKIP_DIRS.iter().map(|dir| (*dir).to_owned()).collect()
293 }
294
295 /// Compares a repository's default branch with the index, whole, and
296 /// queues what differs.
297 async fn reconcile(&self, repo_id: &str) -> Result<()> {
298 let listed: FileList = g1t_kit::call(
299 &self.repos,
300 "list_files",
301 &ListFilesArgs { repo_id: repo_id.to_owned(), git_ref: None, skip_dirs: Self::skip_dirs(), limit: MAX_REPO_FILES },
302 )
303 .await?;
304 let stored: HashMap<String, String> = store::all::<Stored>(
305 &self.db,
306 &Sql { text: "SELECT path, blob FROM files WHERE repo_id = ?".into(), params: vec![p(repo_id)] },
307 )
308 .await?
309 .into_iter()
310 .map(|row| (row.path, row.blob))
311 .collect();
312 let mut changes: Vec<(String, Option<String>)> = Vec::new();
313 let present: HashSet<&str> = listed.files.iter().map(|file| file.path.as_str()).collect();
314 for file in &listed.files {
315 if let Some(hash) = &file.hash
316 && stored.get(&file.path) != Some(hash)
317 {
318 changes.push((file.path.clone(), Some(hash.clone())));
319 }
320 }
321 // Files the branch no longer has; unknowable when the list was cut short.
322 if !listed.truncated {
323 for path in stored.keys().filter(|path| !present.contains(path.as_str())) {
324 changes.push((path.clone(), None));
325 }
326 }
327 self.queue_files(repo_id, &changes).await?;
328 store::run_all(
329 &self.db,
330 vec![store::prepare(
331 &self.db,
332 "UPDATE repos SET state = ?, head = ? WHERE repo_id = ?",
333 vec![p(if listed.truncated { "partial" } else { "indexing" }), Param::opt(listed.commit.as_deref()), p(repo_id)],
334 )?],
335 )
336 .await?;
337 self.drain(repo_id).await
338 }
339
340 /// Reads the next files waiting for a repository, within this job's
341 /// allowance, and queues another job for the rest.
342 pub async fn drain(&self, repo_id: &str) -> Result<()> {
343 let waiting = store::all::<Waiting>(
344 &self.db,
345 &Sql {
346 text: "SELECT path, blob FROM pending WHERE repo_id = ? ORDER BY path LIMIT ?".into(),
347 params: vec![p(repo_id), Param::Int(DRAIN_FILES as i64 + 1)],
348 },
349 )
350 .await?;
351 let more = waiting.len() > DRAIN_FILES;
352 let mut statements = Vec::new();
353 let mut done: Vec<&Waiting> = Vec::new();
354 let mut to_read: Vec<&Waiting> = Vec::new();
355 for file in waiting.iter().take(DRAIN_FILES) {
356 match &file.blob {
357 None => {
358 statements.extend(self.remove_file(repo_id, &file.path)?);
359 done.push(file);
360 }
361 Some(blob) => match rules::skip_path(&file.path) {
362 Some(reason) => {
363 statements.extend(self.skip_file(repo_id, &file.path, blob, 0, reason)?);
364 done.push(file);
365 }
366 None => to_read.push(file),
367 },
368 }
369 }
370 let mut bytes = 0u64;
371 let mut stopped = false;
372 for group in to_read.chunks(READ_GROUP) {
373 if bytes >= DRAIN_BYTES {
374 stopped = true;
375 break;
376 }
377 let hashes: Vec<String> = group.iter().filter_map(|file| file.blob.clone()).collect();
378 let texts: Vec<BlobText> = g1t_kit::call(
379 &self.repos,
380 "read_blobs",
381 &ReadBlobsArgs { repo_id: repo_id.to_owned(), hashes, max_bytes: MAX_FILE_BYTES },
382 )
383 .await?;
384 let by_hash: HashMap<&str, &BlobText> = texts.iter().map(|text| (text.hash.as_str(), text)).collect();
385 for file in group {
386 let blob = file.blob.as_deref().unwrap_or_default();
387 let read = by_hash.get(blob);
388 let size = read.map_or(0, |read| read.size);
389 match read.and_then(|read| read.text.as_deref()) {
390 None => {
391 let reason = if size > u64::from(MAX_FILE_BYTES) { "too large" } else { "binary" };
392 statements.extend(self.skip_file(repo_id, &file.path, blob, size, reason)?);
393 }
394 Some(text) => match rules::skip_text(text) {
395 Some(reason) => statements.extend(self.skip_file(repo_id, &file.path, blob, size, reason)?),
396 None => {
397 statements.extend(self.index_file(repo_id, &file.path, blob, text)?);
398 bytes += text.len() as u64;
399 }
400 },
401 }
402 done.push(file);
403 }
404 }
405 // Each queued file leaves the queue only if no newer push replaced it.
406 for file in &done {
407 statements.push(store::prepare(
408 &self.db,
409 "DELETE FROM pending WHERE repo_id = ? AND path = ? AND blob IS ?",
410 vec![p(repo_id), p(&file.path), Param::opt(file.blob.as_deref())],
411 )?);
412 }
413 store::run_all(&self.db, statements).await?;
414 if more || stopped {
415 return self.enqueue(vec![Job::Drain { repo_id: repo_id.to_owned() }]).await;
416 }
417 self.finish(repo_id).await
418 }
419
420 /// Records a repository's size and language once its queue is empty.
421 async fn finish(&self, repo_id: &str) -> Result<()> {
422 let data: Vec<&str> = rules::DATA_LANGUAGES.to_vec();
423 let state = store::all::<RepoState>(
424 &self.db,
425 &Sql { text: "SELECT state FROM repos WHERE repo_id = ?".into(), params: vec![p(repo_id)] },
426 )
427 .await?;
428 let partial = state.first().is_some_and(|row| row.state == "partial");
429 store::run_all(
430 &self.db,
431 vec![store::prepare(
432 &self.db,
433 "UPDATE repos SET
434 files = (SELECT count(*) FROM files WHERE repo_id = ?1 AND skipped IS NULL),
435 bytes = (SELECT COALESCE(sum(bytes), 0) FROM files WHERE repo_id = ?1 AND skipped IS NULL),
436 language = (SELECT language FROM files WHERE repo_id = ?1 AND skipped IS NULL AND language IS NOT NULL
437 AND language NOT IN (SELECT value FROM json_each(?2))
438 GROUP BY language ORDER BY sum(bytes) DESC LIMIT 1),
439 state = ?3
440 WHERE repo_id = ?1",
441 vec![p(repo_id), Param::Text(serde_json::to_string(&data)?), p(if partial { "partial" } else { "done" })],
442 )?],
443 )
444 .await
445 }
446
447 fn remove_file(&self, repo_id: &str, path: &str) -> Result<Vec<worker::D1PreparedStatement>> {
448 Ok(vec![
449 store::prepare(
450 &self.db,
451 "DELETE FROM chunks WHERE fid = (SELECT fid FROM files WHERE repo_id = ? AND path = ?)",
452 vec![p(repo_id), p(path)],
453 )?,
454 store::prepare(&self.db, "DELETE FROM files WHERE repo_id = ? AND path = ?", vec![p(repo_id), p(path)])?,
455 ])
456 }
457
458 fn upsert_file(&self, repo_id: &str, path: &str, blob: &str, bytes: u64, skipped: Option<&str>) -> Result<worker::D1PreparedStatement> {
459 store::prepare(
460 &self.db,
461 "INSERT INTO files (repo_id, path, blob, language, bytes, skipped) VALUES (?, ?, ?, ?, ?, ?)
462 ON CONFLICT (repo_id, path) DO UPDATE SET blob = excluded.blob, language = excluded.language,
463 bytes = excluded.bytes, skipped = excluded.skipped",
464 vec![p(repo_id), p(path), p(blob), Param::opt(rules::language(path)), Param::Int(bytes as i64), Param::opt(skipped)],
465 )
466 }
467
468 /// Records a file as not indexed, and why, so it is not read again
469 /// until its blob changes.
470 fn skip_file(&self, repo_id: &str, path: &str, blob: &str, bytes: u64, reason: &str) -> Result<Vec<worker::D1PreparedStatement>> {
471 Ok(vec![
472 store::prepare(
473 &self.db,
474 "DELETE FROM chunks WHERE fid = (SELECT fid FROM files WHERE repo_id = ? AND path = ?)",
475 vec![p(repo_id), p(path)],
476 )?,
477 self.upsert_file(repo_id, path, blob, bytes, Some(reason))?,
478 ])
479 }
480
481 fn index_file(&self, repo_id: &str, path: &str, blob: &str, text: &str) -> Result<Vec<worker::D1PreparedStatement>> {
482 let mut statements = vec![
483 self.upsert_file(repo_id, path, blob, text.len() as u64, None)?,
484 store::prepare(
485 &self.db,
486 "DELETE FROM chunks WHERE fid = (SELECT fid FROM files WHERE repo_id = ? AND path = ?)",
487 vec![p(repo_id), p(path)],
488 )?,
489 ];
490 let pieces = rules::chunks(text);
491 for group in pieces.chunks(PIECES_PER_STATEMENT) {
492 let rows = vec!["((SELECT fid FROM files WHERE repo_id = ? AND path = ?), ?, ?, ?)"; group.len()].join(", ");
493 let mut params = Vec::with_capacity(group.len() * 5);
494 for (start, content) in group {
495 params.extend([p(repo_id), p(path), Param::Int(i64::from(*start)), p(path), p(content)]);
496 }
497 statements.push(store::prepare(
498 &self.db,
499 &format!("INSERT INTO chunks (fid, start_line, path, content) VALUES {rows}"),
500 params,
501 )?);
502 }
503 Ok(statements)
504 }
505
506 // ---- Issues and pull requests ------------------------------------------
507
508 fn issue_row(&self, repo_id: &str, issue: &Issue) -> Result<worker::D1PreparedStatement> {
509 let status = match issue.state {
510 State::Open => "open",
511 State::Closed => issue.reason.map_or("completed", |reason| reason.as_str()),
512 };
513 self.item_row(
514 repo_id,
515 "issue",
516 issue.number,
517 &issue.title,
518 &issue.body,
519 if issue.state == State::Open { "open" } else { "closed" },
520 status,
521 (&issue.author.username, issue.requested_by.as_ref().map(|user| user.username.as_str())),
522 &issue.labels,
523 &issue.created_at,
524 &issue.updated_at,
525 )
526 }
527
528 fn pull_row(&self, repo_id: &str, pull: &Pull, labels: &[String]) -> Result<worker::D1PreparedStatement> {
529 self.item_row(
530 repo_id,
531 "pull",
532 pull.number,
533 &pull.title,
534 pull.body.as_deref().unwrap_or_default(),
535 if pull.status.is_active() { "open" } else { "closed" },
536 pull.status.as_str(),
537 (&pull.author.username, pull.requested_by.as_ref().map(|user| user.username.as_str())),
538 labels,
539 &pull.created_at,
540 &pull.updated_at,
541 )
542 }
543
544 #[allow(clippy::too_many_arguments)]
545 fn item_row(
546 &self,
547 repo_id: &str,
548 kind: &str,
549 number: u32,
550 title: &str,
551 body: &str,
552 state: &str,
553 status: &str,
554 // Who opened it, and for g1t's work, who asked for it.
555 (author, requested_by): (&str, Option<&str>),
556 labels: &[String],
557 created_at: &str,
558 updated_at: &str,
559 ) -> Result<worker::D1PreparedStatement> {
560 store::prepare(
561 &self.db,
562 "INSERT INTO items (repo_id, kind, number, title, body, state, status, author, requested_by, labels, created_at, updated_at)
563 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
564 ON CONFLICT (repo_id, kind, number) DO UPDATE SET title = excluded.title, body = excluded.body,
565 state = excluded.state, status = excluded.status, author = excluded.author,
566 requested_by = excluded.requested_by, labels = excluded.labels, updated_at = excluded.updated_at",
567 vec![
568 p(repo_id),
569 p(kind),
570 Param::Int(i64::from(number)),
571 p(title),
572 p(&clip(body, BODY_CHARS)),
573 p(state),
574 p(status),
575 p(&author.to_lowercase()),
576 requested_by.map_or(Param::Null, |name| p(&name.to_lowercase())),
577 p(&labels_of(labels)),
578 p(created_at),
579 p(updated_at),
580 ],
581 )
582 }
583
584 /// Indexes one issue or pull request as it is now.
585 async fn index_item(&self, pull: bool, repo_id: &str, number: u32) -> Result<()> {
586 let Some(path) = self.path_of(repo_id).await? else {
587 return Ok(());
588 };
589 let Some(actor) = self.actor(&path.namespace).await? else {
590 return Ok(());
591 };
592 // An issue on a repository not yet indexed brings the repository in.
593 let known = store::count(
594 &self.db,
595 &Sql { text: "SELECT count(*) AS n FROM repos WHERE repo_id = ?".into(), params: vec![p(repo_id)] },
596 )
597 .await?;
598 if known == 0 && self.refresh_repo(repo_id).await?.is_none() {
599 return Ok(());
600 }
601 let view = ViewArgs { repo: path, number, viewer: Some(actor), after_seq: 0 };
602 let statement = if pull {
603 let found: Outcome<PullDetail> = g1t_kit::call(&self.work, "get_pull", &view).await?;
604 let Outcome::Ok(detail) = found else { return Ok(()) };
605 let labels = detail.issue.as_ref().map(|issue| issue.labels.clone()).unwrap_or_default();
606 self.pull_row(repo_id, &detail.pull, &labels)?
607 } else {
608 let found: Outcome<IssueDetail> = g1t_kit::call(&self.work, "get_issue", &view).await?;
609 let Outcome::Ok(detail) = found else { return Ok(()) };
610 self.issue_row(repo_id, &detail.issue)?
611 };
612 store::run_all(&self.db, vec![statement]).await
613 }
614
615 /// Indexes a repository's issues and pull requests, newest first, as
616 /// many as one list returns.
617 async fn index_items(&self, repo_id: &str, path: &RepoPath) -> Result<()> {
618 let Some(actor) = self.actor(&path.namespace).await? else {
619 return Ok(());
620 };
621 let issues: Outcome<Vec<Issue>> = g1t_kit::call(
622 &self.work,
623 "list_issues",
624 &ListIssuesArgs { repo: path.clone(), viewer: Some(actor.clone()), state: None, label: None },
625 )
626 .await?;
627 let pulls: Outcome<Vec<Pull>> = g1t_kit::call(
628 &self.work,
629 "list_pulls",
630 &ListPullsArgs { repo: path.clone(), viewer: Some(actor), state: None },
631 )
632 .await?;
633 let issues = match issues {
634 Outcome::Ok(issues) => issues,
635 Outcome::Fail(_) => Vec::new(),
636 };
637 let labels: HashMap<u32, &Vec<String>> = issues.iter().map(|issue| (issue.number, &issue.labels)).collect();
638 let mut statements = Vec::new();
639 for issue in &issues {
640 statements.push(self.issue_row(repo_id, issue)?);
641 }
642 if let Outcome::Ok(pulls) = pulls {
643 for pull in &pulls {
644 let labels = pull.issue.and_then(|number| labels.get(&number)).map(|l| l.as_slice()).unwrap_or_default();
645 statements.push(self.pull_row(repo_id, pull, labels)?);
646 }
647 }
648 store::run_all(&self.db, statements).await
649 }
650
651 // ---- People and workspaces ---------------------------------------------
652
653 fn person_row(&self, kind: &str, reference: &str, entry: &DirectoryEntry) -> Result<worker::D1PreparedStatement> {
654 store::prepare(
655 &self.db,
656 "INSERT INTO people (kind, ref, slug, name, bio, avatar, created_at) VALUES (?, ?, ?, ?, ?, ?, ?)
657 ON CONFLICT (kind, ref) DO UPDATE SET slug = excluded.slug, name = excluded.name, bio = excluded.bio,
658 avatar = excluded.avatar",
659 vec![
660 p(kind),
661 p(reference),
662 p(&entry.slug.to_lowercase()),
663 Param::opt(entry.name.as_deref()),
664 Param::opt(entry.bio.as_deref()),
665 Param::opt(entry.avatar.as_deref()),
666 p(&entry.created_at),
667 ],
668 )
669 }
670
671 async fn index_user(&self, username: &str) -> Result<()> {
672 let username = username.to_lowercase();
673 let profile: Option<Profile> = g1t_kit::call(&self.identity, "profile", &UsernameArgs { username: username.clone() }).await?;
674 let statement = match profile {
675 Some(profile) => self.person_row(
676 "user",
677 &username,
678 &DirectoryEntry {
679 id: String::new(),
680 slug: profile.username,
681 name: profile.name,
682 bio: profile.bio,
683 avatar: profile.avatar,
684 created_at: profile.created_at,
685 },
686 )?,
687 None => store::prepare(&self.db, "DELETE FROM people WHERE kind = 'user' AND ref = ?", vec![p(&username)])?,
688 };
689 store::run_all(&self.db, vec![statement]).await
690 }
691
692 async fn index_workspace(&self, slug: &str) -> Result<()> {
693 let workspace: Option<Workspace> = g1t_kit::call(&self.identity, "get_workspace", &SlugArgs { slug: slug.to_lowercase() }).await?;
694 let Some(workspace) = workspace else {
695 return Ok(());
696 };
697 let entry = DirectoryEntry {
698 id: workspace.id.clone(),
699 slug: workspace.slug,
700 name: Some(workspace.name),
701 bio: workspace.description,
702 avatar: workspace.avatar,
703 created_at: workspace.created_at,
704 };
705 store::run_all(&self.db, vec![self.person_row("workspace", &workspace.id, &entry)?]).await
706 }
707
708 // ---- Events and jobs -----------------------------------------------------
709
710 /// Starts the backfill the first time the service hears anything, so a
711 /// new deployment fills itself without anyone asking.
712 pub async fn ensure_backfill(&self) -> Result<()> {
713 #[derive(Deserialize)]
714 struct Started {
715 #[allow(dead_code)]
716 key: String,
717 }
718 let started = store::all::<Started>(
719 &self.db,
720 &Sql {
721 text: "INSERT INTO meta (key, value) VALUES ('backfill', ?) ON CONFLICT (key) DO NOTHING RETURNING key".into(),
722 params: vec![p(&now())],
723 },
724 )
725 .await?;
726 if !started.is_empty() {
727 self.enqueue(vec![
728 Job::Backfill { after: None },
729 Job::Directory { kind: "user".into(), after: None },
730 Job::Directory { kind: "workspace".into(), after: None },
731 ])
732 .await?;
733 }
734 Ok(())
735 }
736
737 pub async fn on_event(&self, event: &Event) -> Result<()> {
738 let data = &event.data;
739 match event.kind.as_str() {
740 "git.push" => {
741 let Ok(push) = serde_json::from_value::<PushEvent>(data.clone()) else { return Ok(()) };
742 if push.default_branch {
743 self.enqueue(vec![Job::Push { repo_id: push.repo_id, before: push.before, after: push.after }]).await?;
744 }
745 }
746 // A transfer changes the repository's path; refreshing reads it
747 // again by id, issues and pull requests with it.
748 // Archiving changes nothing searched, but the details are read again.
749 "repo.created"
750 | "repo.updated"
751 | "repo.visibility_changed"
752 | "repo.renamed"
753 | "repo.transferred"
754 | "repo.archived"
755 | "repo.unarchived" => {
756 if let Ok(repo) = serde_json::from_value::<RepoEvent>(data.clone()) {
757 self.refresh_repo(&repo.repo_id).await?;
758 }
759 }
760 // A deleted repository is forgotten at once (repos hides it, so a
761 // refresh would forget it too), and purged for good later.
762 "repo.deleted" | "repo.purged" => {
763 if let Ok(repo) = serde_json::from_value::<RepoEvent>(data.clone()) {
764 self.purge(&repo.repo_id).await?;
765 }
766 }
767 // A restored one is indexed again, whole: its code, issues and
768 // pull requests.
769 "repo.restored" => {
770 if let Ok(repo) = serde_json::from_value::<RepoEvent>(data.clone()) {
771 self.enqueue(vec![Job::Repo { repo_id: repo.repo_id }]).await?;
772 }
773 }
774 "issue.opened" | "issue.updated" | "issue.closed" | "issue.reopened" | "issue.assigned" => {
775 if let Ok(item) = serde_json::from_value::<NumberEvent>(data.clone()) {
776 self.index_item(false, &item.repo_id, item.number).await?;
777 }
778 }
779 "pull.opened" | "pull.ready" | "pull.updated" | "pull.closed" | "pull.merged" => {
780 if let Ok(item) = serde_json::from_value::<NumberEvent>(data.clone()) {
781 self.index_item(true, &item.repo_id, item.number).await?;
782 }
783 }
784 "user.updated" => {
785 if let Ok(user) = serde_json::from_value::<UserEvent>(data.clone()) {
786 self.index_user(&user.username).await?;
787 }
788 }
789 "workspace.updated" => {
790 if let Ok(workspace) = serde_json::from_value::<WorkspaceEvent>(data.clone()) {
791 self.index_workspace(&workspace.slug).await?;
792 }
793 }
794 "workspace.deleted" => {
795 g1t_kit::deleted::on_event(
796 &self.db,
797 event,
798 &["DELETE FROM people WHERE kind = 'workspace' AND slug = ?1"],
799 )
800 .await?;
801 }
802 "workspace.renamed" => {
803 g1t_kit::rename::on_event(
804 &self.env,
805 &self.db,
806 event,
807 &[
808 "UPDATE repos SET namespace = ?1 WHERE namespace = ?2",
809 "UPDATE people SET slug = ?1 WHERE kind = 'workspace' AND slug = ?2",
810 ],
811 )
812 .await?;
813 }
814 _ => {}
815 }
816 Ok(())
817 }
818
819 pub async fn run_job(&self, job: Job) -> Result<()> {
820 match job {
821 Job::Backfill { after } => {
822 let page: IdPage = g1t_kit::call(&self.repos, "all_ids", &AllIdsArgs { after, limit: BACKFILL_PAGE }).await?;
823 let mut jobs: Vec<Job> = page.ids.into_iter().map(|repo_id| Job::Repo { repo_id }).collect();
824 if let Some(next) = page.next {
825 jobs.push(Job::Backfill { after: Some(next) });
826 }
827 self.enqueue(jobs).await
828 }
829 Job::Repo { repo_id } => {
830 let Some(tree) = self.refresh_repo(&repo_id).await? else { return Ok(()) };
831 let path = RepoPath { namespace: tree.repo.namespace.clone(), name: tree.repo.name.clone() };
832 if let Some(head) = &tree.head {
833 store::run_all(
834 &self.db,
835 vec![store::prepare(
836 &self.db,
837 "UPDATE repos SET pushed_at = COALESCE(pushed_at, ?) WHERE repo_id = ?",
838 vec![p(&head.authored_at), p(&repo_id)],
839 )?],
840 )
841 .await?;
842 }
843 self.index_items(&repo_id, &path).await?;
844 self.reconcile(&repo_id).await
845 }
846 Job::Push { repo_id, before, after } => {
847 if self.refresh_repo(&repo_id).await?.is_none() {
848 return Ok(());
849 }
850 store::run_all(
851 &self.db,
852 vec![store::prepare(&self.db, "UPDATE repos SET pushed_at = ? WHERE repo_id = ?", vec![p(&now()), p(&repo_id)])?],
853 )
854 .await?;
855 let Some(before) = before.filter(|before| !is_null_commit(before)) else {
856 return self.reconcile(&repo_id).await;
857 };
858 let changed: FileList = g1t_kit::call(
859 &self.repos,
860 "changed_files",
861 &ChangedFilesArgs {
862 repo_id: repo_id.clone(),
863 base: Some(before),
864 head: after.clone(),
865 skip_dirs: Self::skip_dirs(),
866 limit: MAX_PUSH_FILES,
867 },
868 )
869 .await?;
870 // A push too large to list, or one whose base is unknown
871 // (a force push over history the store no longer has): compare whole.
872 if changed.truncated || changed.commit.is_none() {
873 return self.reconcile(&repo_id).await;
874 }
875 let files: Vec<(String, Option<String>)> =
876 changed.files.into_iter().map(|file| (file.path, file.hash)).collect();
877 self.queue_files(&repo_id, &files).await?;
878 store::run_all(
879 &self.db,
880 vec![store::prepare(&self.db, "UPDATE repos SET head = ? WHERE repo_id = ?", vec![p(&after), p(&repo_id)])?],
881 )
882 .await?;
883 self.drain(&repo_id).await
884 }
885 Job::Drain { repo_id } => self.drain(&repo_id).await,
886 Job::Directory { kind, after } => {
887 let page: DirectoryPage = g1t_kit::call(
888 &self.identity,
889 "directory",
890 &DirectoryArgs { kind: kind.clone(), after, limit: DIRECTORY_PAGE },
891 )
892 .await?;
893 let mut statements = Vec::new();
894 for entry in &page.entries {
895 let reference = if kind == "workspace" { entry.id.clone() } else { entry.slug.to_lowercase() };
896 statements.push(self.person_row(&kind, &reference, entry)?);
897 }
898 store::run_all(&self.db, statements).await?;
899 if let Some(next) = page.next {
900 self.enqueue(vec![Job::Directory { kind, after: Some(next) }]).await?;
901 }
902 Ok(())
903 }
904 }
905 }
906}
907
908#[cfg(test)]
909mod tests {
910 use g1t_contracts::work::PullStatus;
911
912 use super::*;
913
914 #[test]
915 fn jobs_are_sent_tagged() {
916 let job = Job::Push { repo_id: "rep_1".into(), before: None, after: "abc".into() };
917 let value = serde_json::to_value(&job).unwrap();
918 assert_eq!(value["type"], "push");
919 assert_eq!(serde_json::from_value::<Job>(value).unwrap(), job);
920 let backfill: Job = serde_json::from_str(r#"{"type":"backfill"}"#).unwrap();
921 assert_eq!(backfill, Job::Backfill { after: None });
922 }
923
924 #[test]
925 fn a_push_from_nothing_is_a_whole_comparison() {
926 assert!(is_null_commit("0000000000000000000000000000000000000000"));
927 assert!(is_null_commit(""));
928 assert!(!is_null_commit("0a1b"));
929 }
930
931 #[test]
932 fn labels_are_kept_between_bars() {
933 assert_eq!(labels_of(&["Bug".into(), "good first issue".into()]), "|bug|good first issue|");
934 assert_eq!(labels_of(&[]), "|");
935 }
936
937 #[test]
938 fn pull_statuses_have_names() {
939 assert_eq!(PullStatus::Merged.as_str(), "merged");
940 }
941}