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