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/index.rs

939 lines37,826 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,
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,
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 author: &str,
555 labels: &[String],
556 created_at: &str,
557 updated_at: &str,
558 ) -> Result<worker::D1PreparedStatement> {
559 store::prepare(
560 &self.db,
561 "INSERT INTO items (repo_id, kind, number, title, body, state, status, author, labels, created_at, updated_at)
562 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)
563 ON CONFLICT (repo_id, kind, number) DO UPDATE SET title = excluded.title, body = excluded.body,
564 state = excluded.state, status = excluded.status, author = excluded.author, labels = excluded.labels,
565 updated_at = excluded.updated_at",
566 vec![
567 p(repo_id),
568 p(kind),
569 Param::Int(i64::from(number)),
570 p(title),
571 p(&clip(body, BODY_CHARS)),
572 p(state),
573 p(status),
574 p(&author.to_lowercase()),
575 p(&labels_of(labels)),
576 p(created_at),
577 p(updated_at),
578 ],
579 )
580 }
581
582 /// Indexes one issue or pull request as it is now.
583 async fn index_item(&self, pull: bool, repo_id: &str, number: u32) -> Result<()> {
584 let Some(path) = self.path_of(repo_id).await? else {
585 return Ok(());
586 };
587 let Some(actor) = self.actor(&path.namespace).await? else {
588 return Ok(());
589 };
590 // An issue on a repository not yet indexed brings the repository in.
591 let known = store::count(
592 &self.db,
593 &Sql { text: "SELECT count(*) AS n FROM repos WHERE repo_id = ?".into(), params: vec![p(repo_id)] },
594 )
595 .await?;
596 if known == 0 && self.refresh_repo(repo_id).await?.is_none() {
597 return Ok(());
598 }
599 let view = ViewArgs { repo: path, number, viewer: Some(actor), after_seq: 0 };
600 let statement = if pull {
601 let found: Outcome<PullDetail> = g1t_kit::call(&self.work, "get_pull", &view).await?;
602 let Outcome::Ok(detail) = found else { return Ok(()) };
603 let labels = detail.issue.as_ref().map(|issue| issue.labels.clone()).unwrap_or_default();
604 self.pull_row(repo_id, &detail.pull, &labels)?
605 } else {
606 let found: Outcome<IssueDetail> = g1t_kit::call(&self.work, "get_issue", &view).await?;
607 let Outcome::Ok(detail) = found else { return Ok(()) };
608 self.issue_row(repo_id, &detail.issue)?
609 };
610 store::run_all(&self.db, vec![statement]).await
611 }
612
613 /// Indexes a repository's issues and pull requests, newest first, as
614 /// many as one list returns.
615 async fn index_items(&self, repo_id: &str, path: &RepoPath) -> Result<()> {
616 let Some(actor) = self.actor(&path.namespace).await? else {
617 return Ok(());
618 };
619 let issues: Outcome<Vec<Issue>> = g1t_kit::call(
620 &self.work,
621 "list_issues",
622 &ListIssuesArgs { repo: path.clone(), viewer: Some(actor.clone()), state: None, label: None },
623 )
624 .await?;
625 let pulls: Outcome<Vec<Pull>> = g1t_kit::call(
626 &self.work,
627 "list_pulls",
628 &ListPullsArgs { repo: path.clone(), viewer: Some(actor), state: None },
629 )
630 .await?;
631 let issues = match issues {
632 Outcome::Ok(issues) => issues,
633 Outcome::Fail(_) => Vec::new(),
634 };
635 let labels: HashMap<u32, &Vec<String>> = issues.iter().map(|issue| (issue.number, &issue.labels)).collect();
636 let mut statements = Vec::new();
637 for issue in &issues {
638 statements.push(self.issue_row(repo_id, issue)?);
639 }
640 if let Outcome::Ok(pulls) = pulls {
641 for pull in &pulls {
642 let labels = pull.issue.and_then(|number| labels.get(&number)).map(|l| l.as_slice()).unwrap_or_default();
643 statements.push(self.pull_row(repo_id, pull, labels)?);
644 }
645 }
646 store::run_all(&self.db, statements).await
647 }
648
649 // ---- People and workspaces ---------------------------------------------
650
651 fn person_row(&self, kind: &str, reference: &str, entry: &DirectoryEntry) -> Result<worker::D1PreparedStatement> {
652 store::prepare(
653 &self.db,
654 "INSERT INTO people (kind, ref, slug, name, bio, avatar, created_at) VALUES (?, ?, ?, ?, ?, ?, ?)
655 ON CONFLICT (kind, ref) DO UPDATE SET slug = excluded.slug, name = excluded.name, bio = excluded.bio,
656 avatar = excluded.avatar",
657 vec![
658 p(kind),
659 p(reference),
660 p(&entry.slug.to_lowercase()),
661 Param::opt(entry.name.as_deref()),
662 Param::opt(entry.bio.as_deref()),
663 Param::opt(entry.avatar.as_deref()),
664 p(&entry.created_at),
665 ],
666 )
667 }
668
669 async fn index_user(&self, username: &str) -> Result<()> {
670 let username = username.to_lowercase();
671 let profile: Option<Profile> = g1t_kit::call(&self.identity, "profile", &UsernameArgs { username: username.clone() }).await?;
672 let statement = match profile {
673 Some(profile) => self.person_row(
674 "user",
675 &username,
676 &DirectoryEntry {
677 id: String::new(),
678 slug: profile.username,
679 name: profile.name,
680 bio: profile.bio,
681 avatar: profile.avatar,
682 created_at: profile.created_at,
683 },
684 )?,
685 None => store::prepare(&self.db, "DELETE FROM people WHERE kind = 'user' AND ref = ?", vec![p(&username)])?,
686 };
687 store::run_all(&self.db, vec![statement]).await
688 }
689
690 async fn index_workspace(&self, slug: &str) -> Result<()> {
691 let workspace: Option<Workspace> = g1t_kit::call(&self.identity, "get_workspace", &SlugArgs { slug: slug.to_lowercase() }).await?;
692 let Some(workspace) = workspace else {
693 return Ok(());
694 };
695 let entry = DirectoryEntry {
696 id: workspace.id.clone(),
697 slug: workspace.slug,
698 name: Some(workspace.name),
699 bio: workspace.description,
700 avatar: workspace.avatar,
701 created_at: workspace.created_at,
702 };
703 store::run_all(&self.db, vec![self.person_row("workspace", &workspace.id, &entry)?]).await
704 }
705
706 // ---- Events and jobs -----------------------------------------------------
707
708 /// Starts the backfill the first time the service hears anything, so a
709 /// new deployment fills itself without anyone asking.
710 pub async fn ensure_backfill(&self) -> Result<()> {
711 #[derive(Deserialize)]
712 struct Started {
713 #[allow(dead_code)]
714 key: String,
715 }
716 let started = store::all::<Started>(
717 &self.db,
718 &Sql {
719 text: "INSERT INTO meta (key, value) VALUES ('backfill', ?) ON CONFLICT (key) DO NOTHING RETURNING key".into(),
720 params: vec![p(&now())],
721 },
722 )
723 .await?;
724 if !started.is_empty() {
725 self.enqueue(vec![
726 Job::Backfill { after: None },
727 Job::Directory { kind: "user".into(), after: None },
728 Job::Directory { kind: "workspace".into(), after: None },
729 ])
730 .await?;
731 }
732 Ok(())
733 }
734
735 pub async fn on_event(&self, event: &Event) -> Result<()> {
736 let data = &event.data;
737 match event.kind.as_str() {
738 "git.push" => {
739 let Ok(push) = serde_json::from_value::<PushEvent>(data.clone()) else { return Ok(()) };
740 if push.default_branch {
741 self.enqueue(vec![Job::Push { repo_id: push.repo_id, before: push.before, after: push.after }]).await?;
742 }
743 }
744 // A transfer changes the repository's path; refreshing reads it
745 // again by id, issues and pull requests with it.
746 // Archiving changes nothing searched, but the details are read again.
747 "repo.created"
748 | "repo.updated"
749 | "repo.visibility_changed"
750 | "repo.renamed"
751 | "repo.transferred"
752 | "repo.archived"
753 | "repo.unarchived" => {
754 if let Ok(repo) = serde_json::from_value::<RepoEvent>(data.clone()) {
755 self.refresh_repo(&repo.repo_id).await?;
756 }
757 }
758 // A deleted repository is forgotten at once (repos hides it, so a
759 // refresh would forget it too), and purged for good later.
760 "repo.deleted" | "repo.purged" => {
761 if let Ok(repo) = serde_json::from_value::<RepoEvent>(data.clone()) {
762 self.purge(&repo.repo_id).await?;
763 }
764 }
765 // A restored one is indexed again, whole: its code, issues and
766 // pull requests.
767 "repo.restored" => {
768 if let Ok(repo) = serde_json::from_value::<RepoEvent>(data.clone()) {
769 self.enqueue(vec![Job::Repo { repo_id: repo.repo_id }]).await?;
770 }
771 }
772 "issue.opened" | "issue.updated" | "issue.closed" | "issue.reopened" | "issue.assigned" => {
773 if let Ok(item) = serde_json::from_value::<NumberEvent>(data.clone()) {
774 self.index_item(false, &item.repo_id, item.number).await?;
775 }
776 }
777 "pull.opened" | "pull.ready" | "pull.updated" | "pull.closed" | "pull.merged" => {
778 if let Ok(item) = serde_json::from_value::<NumberEvent>(data.clone()) {
779 self.index_item(true, &item.repo_id, item.number).await?;
780 }
781 }
782 "user.updated" => {
783 if let Ok(user) = serde_json::from_value::<UserEvent>(data.clone()) {
784 self.index_user(&user.username).await?;
785 }
786 }
787 "workspace.updated" => {
788 if let Ok(workspace) = serde_json::from_value::<WorkspaceEvent>(data.clone()) {
789 self.index_workspace(&workspace.slug).await?;
790 }
791 }
792 "workspace.deleted" => {
793 g1t_kit::deleted::on_event(
794 &self.db,
795 event,
796 &["DELETE FROM people WHERE kind = 'workspace' AND slug = ?1"],
797 )
798 .await?;
799 }
800 "workspace.renamed" => {
801 g1t_kit::rename::on_event(
802 &self.env,
803 &self.db,
804 event,
805 &[
806 "UPDATE repos SET namespace = ?1 WHERE namespace = ?2",
807 "UPDATE people SET slug = ?1 WHERE kind = 'workspace' AND slug = ?2",
808 ],
809 )
810 .await?;
811 }
812 _ => {}
813 }
814 Ok(())
815 }
816
817 pub async fn run_job(&self, job: Job) -> Result<()> {
818 match job {
819 Job::Backfill { after } => {
820 let page: IdPage = g1t_kit::call(&self.repos, "all_ids", &AllIdsArgs { after, limit: BACKFILL_PAGE }).await?;
821 let mut jobs: Vec<Job> = page.ids.into_iter().map(|repo_id| Job::Repo { repo_id }).collect();
822 if let Some(next) = page.next {
823 jobs.push(Job::Backfill { after: Some(next) });
824 }
825 self.enqueue(jobs).await
826 }
827 Job::Repo { repo_id } => {
828 let Some(tree) = self.refresh_repo(&repo_id).await? else { return Ok(()) };
829 let path = RepoPath { namespace: tree.repo.namespace.clone(), name: tree.repo.name.clone() };
830 if let Some(head) = &tree.head {
831 store::run_all(
832 &self.db,
833 vec![store::prepare(
834 &self.db,
835 "UPDATE repos SET pushed_at = COALESCE(pushed_at, ?) WHERE repo_id = ?",
836 vec![p(&head.authored_at), p(&repo_id)],
837 )?],
838 )
839 .await?;
840 }
841 self.index_items(&repo_id, &path).await?;
842 self.reconcile(&repo_id).await
843 }
844 Job::Push { repo_id, before, after } => {
845 if self.refresh_repo(&repo_id).await?.is_none() {
846 return Ok(());
847 }
848 store::run_all(
849 &self.db,
850 vec![store::prepare(&self.db, "UPDATE repos SET pushed_at = ? WHERE repo_id = ?", vec![p(&now()), p(&repo_id)])?],
851 )
852 .await?;
853 let Some(before) = before.filter(|before| !is_null_commit(before)) else {
854 return self.reconcile(&repo_id).await;
855 };
856 let changed: FileList = g1t_kit::call(
857 &self.repos,
858 "changed_files",
859 &ChangedFilesArgs {
860 repo_id: repo_id.clone(),
861 base: Some(before),
862 head: after.clone(),
863 skip_dirs: Self::skip_dirs(),
864 limit: MAX_PUSH_FILES,
865 },
866 )
867 .await?;
868 // A push too large to list, or one whose base is unknown
869 // (a force push over history the store no longer has): compare whole.
870 if changed.truncated || changed.commit.is_none() {
871 return self.reconcile(&repo_id).await;
872 }
873 let files: Vec<(String, Option<String>)> =
874 changed.files.into_iter().map(|file| (file.path, file.hash)).collect();
875 self.queue_files(&repo_id, &files).await?;
876 store::run_all(
877 &self.db,
878 vec![store::prepare(&self.db, "UPDATE repos SET head = ? WHERE repo_id = ?", vec![p(&after), p(&repo_id)])?],
879 )
880 .await?;
881 self.drain(&repo_id).await
882 }
883 Job::Drain { repo_id } => self.drain(&repo_id).await,
884 Job::Directory { kind, after } => {
885 let page: DirectoryPage = g1t_kit::call(
886 &self.identity,
887 "directory",
888 &DirectoryArgs { kind: kind.clone(), after, limit: DIRECTORY_PAGE },
889 )
890 .await?;
891 let mut statements = Vec::new();
892 for entry in &page.entries {
893 let reference = if kind == "workspace" { entry.id.clone() } else { entry.slug.to_lowercase() };
894 statements.push(self.person_row(&kind, &reference, entry)?);
895 }
896 store::run_all(&self.db, statements).await?;
897 if let Some(next) = page.next {
898 self.enqueue(vec![Job::Directory { kind, after: Some(next) }]).await?;
899 }
900 Ok(())
901 }
902 }
903 }
904}
905
906#[cfg(test)]
907mod tests {
908 use g1t_contracts::work::PullStatus;
909
910 use super::*;
911
912 #[test]
913 fn jobs_are_sent_tagged() {
914 let job = Job::Push { repo_id: "rep_1".into(), before: None, after: "abc".into() };
915 let value = serde_json::to_value(&job).unwrap();
916 assert_eq!(value["type"], "push");
917 assert_eq!(serde_json::from_value::<Job>(value).unwrap(), job);
918 let backfill: Job = serde_json::from_str(r#"{"type":"backfill"}"#).unwrap();
919 assert_eq!(backfill, Job::Backfill { after: None });
920 }
921
922 #[test]
923 fn a_push_from_nothing_is_a_whole_comparison() {
924 assert!(is_null_commit("0000000000000000000000000000000000000000"));
925 assert!(is_null_commit(""));
926 assert!(!is_null_commit("0a1b"));
927 }
928
929 #[test]
930 fn labels_are_kept_between_bars() {
931 assert_eq!(labels_of(&["Bug".into(), "good first issue".into()]), "|bug|good first issue|");
932 assert_eq!(labels_of(&[]), "|");
933 }
934
935 #[test]
936 fn pull_statuses_have_names() {
937 assert_eq!(PullStatus::Merged.as_str(), "merged");
938 }
939}