Skip to content
994 linesCodeBlameRaw
1//! Repository metadata in D1.
2
3use std::cell::RefCell;
4use std::collections::HashMap;
5
6use g1t_contracts::Viewer;
7use g1t_contracts::access::{self, Capability, RepoRole};
8use g1t_contracts::repos::{Repo, RepoPath};
9use serde::Deserialize;
10use worker::wasm_bindgen::JsValue;
11use worker::{D1Database, Result};
12
13#[derive(Deserialize)]
14pub(crate) struct RepoRow {
15 id: String,
16 namespace: String,
17 name: String,
18 description: Option<String>,
19 is_private: u8,
20 owner_id: String,
21 default_branch: String,
22 fork_of: Option<String>,
23 protected: u8,
24 created_at: String,
25 /// Null only on rows written before the column existed and not yet
26 /// migrated; their key is the one worked out from the path.
27 #[serde(default)]
28 store: Option<String>,
29 /// JSON; absent on rows read before the column existed.
30 #[serde(default)]
31 topics: Option<String>,
32 #[serde(default)]
33 website: Option<String>,
34 #[serde(default)]
35 archived_at: Option<String>,
36 /// JSON `RepoMirror`; absent on rows read before the column existed.
37 #[serde(default)]
38 mirror: Option<String>,
39 #[serde(default)]
40 deleted_at: Option<String>,
41 /// Bumped by everything that changes the repository's refs; see
42 /// [`RefsState`]. Absent on rows read before the column existed.
43 #[serde(default)]
44 refs_version: Option<f64>,
45 #[serde(default)]
46 refs_open_until: Option<f64>,
47 /// A pull request working copy whose git data was removed, and the
48 /// head it had (forks.rs). Absent before the columns existed.
49 #[serde(default)]
50 retired_at: Option<String>,
51 #[serde(default)]
52 retired_head: Option<String>,
53 /// Until when writes wait, and why: a move between namespaces
54 /// (moves.rs). Absent before the columns existed.
55 #[serde(default)]
56 writes_paused_until: Option<f64>,
57 #[serde(default)]
58 writes_paused_for: Option<String>,
59}
60
61thread_local! {
62 /// Repositories whose writes wait, by id: until when, and why. Filled
63 /// whenever a row is read.
64 static PAUSED: RefCell<HashMap<String, (u64, String)>> = RefCell::new(HashMap::new());
65}
66
67/// Records whether writes to the repository with this id wait, as its row says.
68pub fn note_paused(id: &str, until: Option<u64>, reason: Option<&str>) {
69 PAUSED.with(|paused| {
70 let mut paused = paused.borrow_mut();
71 match until {
72 Some(until) => {
73 paused.insert(id.to_owned(), (until, reason.unwrap_or("maintenance").to_owned()));
74 }
75 None => {
76 paused.remove(id);
77 }
78 }
79 });
80}
81
82/// Why writes to the repository with this id wait at `now`, if they do, as
83/// its row last read here said.
84pub fn paused(id: &str, now: u64) -> Option<String> {
85 PAUSED.with(|paused| paused.borrow().get(id).filter(|(until, _)| *until > now).map(|(_, reason)| reason.clone()))
86}
87
88/// Where a repository's refs stand, as its row last said: `version` goes up
89/// with every change g1t makes to them, so an answer that lists them (see
90/// refs_cache.rs) is kept under the version it was made at, and a change
91/// leaves it behind. Until `open_until` (milliseconds) a credential that
92/// can change them is out of g1t's hands, and nothing is kept.
93#[derive(Clone, Copy, Debug, Default, PartialEq, Eq)]
94pub struct RefsState {
95 pub version: u64,
96 pub open_until: u64,
97}
98
99/// The newest [`RefsState`] this isolate has read or written, by
100/// repository id. A version only goes up, so an older read finishing late
101/// never takes a newer one back.
102#[derive(Default)]
103pub struct RefsStates {
104 states: HashMap<String, RefsState>,
105}
106
107impl RefsStates {
108 pub fn note(&mut self, id: &str, state: RefsState) {
109 let kept = self.states.entry(id.to_owned()).or_default();
110 kept.version = kept.version.max(state.version);
111 kept.open_until = kept.open_until.max(state.open_until);
112 }
113
114 pub fn get(&self, id: &str) -> Option<RefsState> {
115 self.states.get(id).copied()
116 }
117}
118
119thread_local! {
120 static REFS: RefCell<RefsStates> = RefCell::new(RefsStates::default());
121}
122
123/// Where the refs of the repository with this id stand, as this isolate
124/// last read them; `None` before the column existed or before its row was
125/// read here.
126pub fn refs_state(id: &str) -> Option<RefsState> {
127 REFS.with(|refs| refs.borrow().get(id))
128}
129
130fn note_refs(id: &str, version: Option<f64>, open_until: Option<f64>) {
131 if let Some(version) = version {
132 let state = RefsState {
133 version: version as u64,
134 open_until: open_until.unwrap_or(0.0) as u64,
135 };
136 REFS.with(|refs| refs.borrow_mut().note(id, state));
137 }
138}
139
140thread_local! {
141 /// Working copies whose git data was removed, by id, with the head
142 /// each had (forks.rs). Filled whenever a row is read.
143 static RETIRED: RefCell<HashMap<String, String>> = RefCell::new(HashMap::new());
144}
145
146/// The head a removed working copy had, if the repository with this id is one.
147pub fn retired(id: &str) -> Option<String> {
148 RETIRED.with(|retired| retired.borrow().get(id).cloned())
149}
150
151/// Records whether the repository with this id is a removed working copy.
152pub fn note_retired(id: &str, head: Option<&str>) {
153 RETIRED.with(|retired| {
154 let mut retired = retired.borrow_mut();
155 match head {
156 Some(head) => {
157 retired.insert(id.to_owned(), head.to_owned());
158 }
159 None => {
160 retired.remove(id);
161 }
162 }
163 });
164}
165
166thread_local! {
167 /// Store keys that differ from the one a repository's path gives: those
168 /// of repositories whose workspace was renamed after they were made.
169 /// Filled whenever a row is read or written, so every `Repo` this
170 /// service holds has its key here. A key changes only when a move
171 /// between namespaces switches it (moves.rs), and every row read
172 /// after that brings the new one, so requests sharing the isolate can
173 /// share the map.
174 static MOVED: RefCell<HashMap<String, String>> = RefCell::new(HashMap::new());
175}
176
177/// How long a fetch may go by a repository's row as it was read a moment
178/// ago: a clone is two or three requests in quick succession, and each
179/// would otherwise read the same row. Short enough that making a repository
180/// private, archiving or deleting it applies within seconds.
181pub const RECENT_MS: u64 = 5_000;
182
183/// Repositories read in the last [`RECENT_MS`], by path. Only rows that
184/// were found are kept, so a repository just made is never missed.
185#[derive(Default)]
186pub struct Recent {
187 rows: HashMap<(String, String), (Repo, u64)>,
188}
189
190impl Recent {
191 fn key(path: &RepoPath) -> (String, String) {
192 (path.namespace.to_lowercase(), path.name.to_lowercase())
193 }
194
195 pub fn get(&self, path: &RepoPath, now: u64) -> Option<Repo> {
196 self.rows
197 .get(&Self::key(path))
198 .filter(|(_, read)| now.saturating_sub(*read) < RECENT_MS)
199 .map(|(repo, _)| repo.clone())
200 }
201
202 pub fn keep(&mut self, path: &RepoPath, repo: &Repo, now: u64) {
203 self.rows.retain(|_, (_, read)| now.saturating_sub(*read) < RECENT_MS);
204 self.rows.insert(Self::key(path), (repo.clone(), now));
205 }
206}
207
208thread_local! {
209 static RECENT: RefCell<Recent> = RefCell::new(Recent::default());
210}
211
212/// The key a repository's path gives: what every repository was stored
213/// under before workspaces could be renamed.
214pub fn path_key(repo: &Repo) -> String {
215 format!("{}--{}", repo.namespace, repo.name)
216}
217
218/// Records where a repository is stored, when its path does not say. A
219/// key changes when the repository moves between namespaces (moves.rs),
220/// so one that is the path's again is forgotten.
221pub fn remember_store(repo: &Repo, store: &str) {
222 MOVED.with(|moved| {
223 let mut moved = moved.borrow_mut();
224 if store != path_key(repo) {
225 moved.insert(repo.id.clone(), store.to_owned());
226 } else {
227 moved.remove(&repo.id);
228 }
229 });
230}
231
232impl From<RepoRow> for Repo {
233 fn from(row: RepoRow) -> Self {
234 let repo = Repo {
235 id: row.id,
236 namespace: row.namespace,
237 name: row.name,
238 description: row.description,
239 is_private: row.is_private != 0,
240 owner_id: row.owner_id,
241 default_branch: row.default_branch,
242 fork_of: row.fork_of,
243 protected: row.protected != 0,
244 created_at: row.created_at,
245 topics: row
246 .topics
247 .as_deref()
248 .and_then(|topics| serde_json::from_str(topics).ok())
249 .unwrap_or_default(),
250 website: row.website,
251 archived_at: row.archived_at,
252 mirror: row.mirror.as_deref().and_then(|mirror| serde_json::from_str(mirror).ok()),
253 };
254 if let Some(store) = &row.store {
255 remember_store(&repo, store);
256 }
257 note_refs(&repo.id, row.refs_version, row.refs_open_until);
258 note_retired(&repo.id, row.retired_at.as_ref().and(row.retired_head.as_deref()));
259 note_paused(&repo.id, row.writes_paused_until.map(|until| until as u64), row.writes_paused_for.as_deref());
260 repo
261 }
262}
263
264/// The key a repo is stored under in the git store.
265pub fn store_key(repo: &Repo) -> String {
266 let key = MOVED
267 .with(|moved| moved.borrow().get(&repo.id).cloned())
268 .unwrap_or_else(|| path_key(repo));
269 // Its interactions with the store are metered for its workspace.
270 crate::meters::note_owner(&key, &repo.namespace);
271 key
272}
273
274/// The viewer's role on `repo` (see `g1t_contracts::access`): ownership of
275/// its workspace, the workspace's base permission, a direct grant, or
276/// Read on a public repository. A pull request's fork is its author's to
277/// write; whoever can read the repository it came from can read it too,
278/// which `Repos::may_read` checks.
279pub fn role(repo: &Repo, viewer: &Viewer) -> Option<RepoRole> {
280 if repo.fork_of.is_some() {
281 let author = viewer.as_ref().is_some_and(|user| user.id == repo.owner_id);
282 return if author {
283 Some(RepoRole::Write)
284 } else if repo.is_private {
285 None
286 } else {
287 Some(RepoRole::Read)
288 };
289 }
290 access::permission(viewer.as_ref(), repo)
291}
292
293/// Whether the viewer may read `repo`, going by the repository alone.
294pub fn can_read(repo: &Repo, viewer: &Viewer) -> bool {
295 role(repo, viewer).is_some()
296}
297
298/// Whether the viewer may push to `repo`: Write or higher, or the author
299/// of a pull request's fork.
300pub fn can_write(repo: &Repo, viewer: &Viewer) -> bool {
301 can(repo, viewer, Capability::Push)
302}
303
304/// Whether the viewer may do `capability` in `repo`. A fork has only its
305/// author's Write.
306pub fn can(repo: &Repo, viewer: &Viewer, capability: Capability) -> bool {
307 if repo.fork_of.is_some() {
308 return role(repo, viewer).is_some_and(|role| access::allows(role, capability))
309 && !access::OWNER_ONLY.contains(&capability);
310 }
311 access::can(viewer.as_ref(), repo, capability)
312}
313
314fn optional(value: &Option<String>) -> JsValue {
315 value.as_deref().map_or(JsValue::NULL, JsValue::from)
316}
317
318pub struct Registry {
319 pub db: D1Database,
320}
321
322impl Registry {
323 pub async fn by_path(&self, path: &RepoPath) -> Result<Option<Repo>> {
324 Ok(self
325 .db
326 .prepare("SELECT * FROM repos WHERE namespace = ? AND name = ? AND deleted_at IS NULL")
327 .bind(&[
328 path.namespace.to_lowercase().into(),
329 path.name.to_lowercase().into(),
330 ])?
331 .first::<RepoRow>(None)
332 .await?
333 .map(Repo::from))
334 }
335
336 /// The repository at `path`, as read in the last few seconds if it was
337 /// (see [`RECENT_MS`]). For fetches only: a push always reads the row.
338 pub async fn by_path_recent(&self, path: &RepoPath) -> Result<Option<Repo>> {
339 let now = g1t_kit::now_ms();
340 if let Some(repo) = RECENT.with(|recent| recent.borrow().get(path, now)) {
341 return Ok(Some(repo));
342 }
343 let found = self.by_path(path).await?;
344 if let Some(repo) = &found {
345 RECENT.with(|recent| recent.borrow_mut().keep(path, repo, now));
346 }
347 Ok(found)
348 }
349
350 /// Its details; who can see it changes with `set_private`.
351 pub async fn update(
352 &self,
353 id: &str,
354 description: Option<&str>,
355 protected: bool,
356 topics: &[String],
357 website: Option<&str>,
358 ) -> Result<()> {
359 self.db
360 .prepare("UPDATE repos SET description = ?, protected = ?, topics = ?, website = ? WHERE id = ?")
361 .bind(&[
362 description.map_or(JsValue::NULL, JsValue::from),
363 u32::from(protected).into(),
364 serde_json::to_string(topics)?.into(),
365 website.map_or(JsValue::NULL, JsValue::from),
366 id.into(),
367 ])?
368 .run()
369 .await?;
370 Ok(())
371 }
372
373 pub async fn by_id(&self, id: &str) -> Result<Option<Repo>> {
374 Ok(self
375 .db
376 .prepare("SELECT * FROM repos WHERE id = ? AND deleted_at IS NULL")
377 .bind(&[id.into()])?
378 .first::<RepoRow>(None)
379 .await?
380 .map(Repo::from))
381 }
382
383 /// The repository at `path`, deleted or not: what holds the name.
384 pub async fn by_path_any(&self, path: &RepoPath) -> Result<Option<(Repo, Option<String>)>> {
385 Ok(self
386 .db
387 .prepare("SELECT * FROM repos WHERE namespace = ? AND name = ?")
388 .bind(&[
389 path.namespace.to_lowercase().into(),
390 path.name.to_lowercase().into(),
391 ])?
392 .first::<RepoRow>(None)
393 .await?
394 .map(|mut row| {
395 let deleted_at = row.deleted_at.take();
396 (Repo::from(row), deleted_at)
397 }))
398 }
399
400 /// Repos the viewer may see, newest first. Excludes pull request forks.
401 /// With `member_only`, only repos in the viewer's own workspaces.
402 pub async fn list(
403 &self,
404 viewer: &Viewer,
405 query: Option<&str>,
406 namespace: Option<&str>,
407 member_only: bool,
408 ) -> Result<Vec<Repo>> {
409 let workspaces: Vec<&str> = viewer
410 .iter()
411 .flat_map(|user| &user.workspaces)
412 .map(|membership| membership.slug.as_str())
413 .collect();
414 // The workspaces whose private repositories the viewer reads all
415 // of (an owner, or a base permission other than none), and the
416 // repositories they were given a role on: see access.rs. A probe
417 // repository in each workspace stands for all of them.
418 let reading: Vec<&str> = viewer
419 .iter()
420 .flat_map(|user| {
421 user.workspaces.iter().filter(move |membership| {
422 let probe = access::RepoRef { id: "", namespace: &membership.slug, private: true };
423 access::granted(user, probe).is_some()
424 })
425 })
426 .map(|membership| membership.slug.as_str())
427 .collect();
428 let mut granted: Vec<&str> = viewer
429 .iter()
430 .flat_map(|user| {
431 let token = user.token.as_deref();
432 user.grants
433 .iter()
434 .filter(move |grant| token.is_none_or(|token| token.covers_repo(&grant.repo_id, &grant.workspace)))
435 })
436 .map(|grant| grant.repo_id.as_str())
437 .collect();
438 // A token's selected repositories, where its owner's
439 // membership reaches them.
440 if let Some(user) = viewer.as_ref()
441 && let Some(reach) = user.token.as_deref().and_then(|token| token.reach.as_ref())
442 && let Some(workspace) = reach.workspace.as_deref()
443 {
444 granted.extend(
445 reach
446 .repo_ids
447 .iter()
448 .filter(|id| access::granted(user, access::RepoRef { id, namespace: workspace, private: true }).is_some())
449 .map(String::as_str),
450 );
451 }
452 let mut params: Vec<JsValue> = vec![
453 serde_json::to_string(&reading)?.into(),
454 serde_json::to_string(&granted)?.into(),
455 ];
456 let private_ok = "(namespace IN (SELECT value FROM json_each(?)) OR id IN (SELECT value FROM json_each(?)))";
457 let mut conditions = vec![
458 "fork_of IS NULL AND deleted_at IS NULL".to_owned(),
459 format!("(is_private = 0 OR {private_ok})"),
460 ];
461 if member_only {
462 conditions.push("namespace IN (SELECT value FROM json_each(?))".to_owned());
463 params.push(serde_json::to_string(&workspaces)?.into());
464 }
465 if let Some(namespace) = namespace {
466 conditions.push("namespace = ?".to_owned());
467 params.push(namespace.to_lowercase().into());
468 }
469 if let Some(query) = query.map(str::trim).filter(|query| !query.is_empty()) {
470 conditions
471 .push("(name LIKE ? ESCAPE '\\' OR description LIKE ? ESCAPE '\\')".to_owned());
472 // LIKE wildcards in the query are matched literally.
473 let escaped: String = query
474 .chars()
475 .flat_map(|c| match c {
476 '\\' | '%' | '_' => vec!['\\', c],
477 _ => vec![c],
478 })
479 .collect();
480 let pattern = format!("%{escaped}%");
481 params.push(pattern.as_str().into());
482 params.push(pattern.into());
483 }
484 let sql = format!(
485 "SELECT * FROM repos WHERE {} ORDER BY created_at DESC, id DESC LIMIT 50",
486 conditions.join(" AND ")
487 );
488 let rows = self
489 .db
490 .prepare(sql)
491 .bind(&params)?
492 .all()
493 .await?
494 .results::<RepoRow>()?;
495 Ok(rows.into_iter().map(Repo::from).collect())
496 }
497
498 /// Of these ids, the repositories (not forks) the viewer may read.
499 pub async fn readable(&self, ids: &[String], viewer: &Viewer) -> Result<Vec<Repo>> {
500 let ids: Vec<&String> = ids.iter().take(g1t_contracts::repos::MAX_READABLE).collect();
501 if ids.is_empty() {
502 return Ok(Vec::new());
503 }
504 // One parameter however many ids: D1 binds at most 100.
505 let rows = self
506 .db
507 .prepare(
508 "SELECT * FROM repos
509 WHERE id IN (SELECT value FROM json_each(?)) AND fork_of IS NULL AND deleted_at IS NULL",
510 )
511 .bind(&[serde_json::to_string(&ids)?.into()])?
512 .all()
513 .await?
514 .results::<RepoRow>()?;
515 Ok(rows
516 .into_iter()
517 .map(Repo::from)
518 .filter(|repo| can_read(repo, viewer))
519 .collect())
520 }
521
522 /// The workspaces in which this account made a public repository.
523 pub async fn public_namespaces(&self, owner_id: &str) -> Result<Vec<String>> {
524 #[derive(Deserialize)]
525 struct Row {
526 namespace: String,
527 }
528 Ok(self
529 .db
530 .prepare(
531 "SELECT DISTINCT namespace FROM repos
532 WHERE owner_id = ? AND is_private = 0 AND fork_of IS NULL AND deleted_at IS NULL
533 ORDER BY namespace",
534 )
535 .bind(&[owner_id.into()])?
536 .all()
537 .await?
538 .results::<Row>()?
539 .into_iter()
540 .map(|row| row.namespace)
541 .collect())
542 }
543
544 /// Repositories that are not forks, by id, a page at a time.
545 /// Repositories that are not forks, with who created each, by id after
546 /// `after`.
547 pub async fn creators_after(&self, after: Option<&str>, limit: u32) -> Result<Vec<g1t_contracts::repos::RepoCreator>> {
548 self.db
549 .prepare(
550 "SELECT id, namespace, name, owner_id FROM repos
551 WHERE fork_of IS NULL AND deleted_at IS NULL AND id > ? ORDER BY id LIMIT ?",
552 )
553 .bind(&[after.unwrap_or("").into(), limit.into()])?
554 .all()
555 .await?
556 .results::<g1t_contracts::repos::RepoCreator>()
557 }
558
559 pub async fn ids_after(&self, after: Option<&str>, limit: u32) -> Result<Vec<String>> {
560 #[derive(Deserialize)]
561 struct Row {
562 id: String,
563 }
564 Ok(self
565 .db
566 .prepare("SELECT id FROM repos WHERE fork_of IS NULL AND deleted_at IS NULL AND id > ? ORDER BY id LIMIT ?")
567 .bind(&[after.unwrap_or("").into(), limit.into()])?
568 .all()
569 .await?
570 .results::<Row>()?
571 .into_iter()
572 .map(|row| row.id)
573 .collect())
574 }
575
576 /// Adds a pushed pack's bytes to what the repository is counted as
577 /// holding: its own, or, for a pull request's working copy, the
578 /// repository it is a copy of, whose storage it is.
579 pub async fn add_stored_bytes(&self, repo: &Repo, bytes: u64) -> Result<()> {
580 if bytes == 0 {
581 return Ok(());
582 }
583 let root = repo.fork_of.as_deref().unwrap_or(&repo.id);
584 self.db
585 .prepare("UPDATE repos SET stored_bytes = stored_bytes + ? WHERE id = ?")
586 .bind(&[(bytes as f64).into(), root.into()])?
587 .run()
588 .await?;
589 Ok(())
590 }
591
592 /// Which of these `namespace/name` paths are private. A working copy
593 /// answers as its repository. Unknown paths are left out.
594 pub async fn visibility(&self, paths: &[String]) -> Result<Vec<g1t_contracts::repos::RepoVisibility>> {
595 let mut out = Vec::new();
596 for path in paths.iter().take(50) {
597 let Some((namespace, name)) = path.split_once('/') else { continue };
598 let Some(repo) = self
599 .by_path(&RepoPath { namespace: namespace.to_owned(), name: name.to_owned() })
600 .await?
601 else {
602 continue;
603 };
604 let is_private = match &repo.fork_of {
605 Some(parent) => self.by_id(parent).await?.map_or(repo.is_private, |parent| parent.is_private),
606 None => repo.is_private,
607 };
608 out.push(g1t_contracts::repos::RepoVisibility { path: path.clone(), is_private });
609 }
610 Ok(out)
611 }
612
613 /// What each workspace's repositories are counted as holding, private
614 /// and public apart. Working copies count toward their repository.
615 pub async fn storage(&self) -> Result<Vec<g1t_contracts::repos::WorkspaceStorage>> {
616 #[derive(Deserialize)]
617 struct Row {
618 namespace: String,
619 private_bytes: Option<f64>,
620 public_bytes: Option<f64>,
621 }
622 Ok(self
623 .db
624 .prepare(
625 "SELECT namespace,
626 SUM(CASE WHEN is_private = 1 THEN stored_bytes ELSE 0 END) AS private_bytes,
627 SUM(CASE WHEN is_private = 0 THEN stored_bytes ELSE 0 END) AS public_bytes
628 FROM repos WHERE fork_of IS NULL AND deleted_at IS NULL AND stored_bytes > 0 GROUP BY namespace",
629 )
630 .all()
631 .await?
632 .results::<Row>()?
633 .into_iter()
634 .map(|row| g1t_contracts::repos::WorkspaceStorage {
635 namespace: row.namespace,
636 private_bytes: row.private_bytes.unwrap_or(0.0) as i64,
637 public_bytes: row.public_bytes.unwrap_or(0.0) as i64,
638 })
639 .collect())
640 }
641
642 /// What one workspace's private repositories are counted as holding.
643 pub async fn private_bytes(&self, namespace: &str) -> Result<i64> {
644 #[derive(Deserialize)]
645 struct Row {
646 bytes: Option<f64>,
647 }
648 Ok(self
649 .db
650 .prepare("SELECT SUM(stored_bytes) AS bytes FROM repos WHERE namespace = ? AND is_private = 1 AND fork_of IS NULL AND deleted_at IS NULL")
651 .bind(&[namespace.into()])?
652 .first::<Row>(None)
653 .await?
654 .and_then(|row| row.bytes)
655 .unwrap_or(0.0) as i64)
656 }
657
658 /// Forgets a repository that could not be filled.
659 pub async fn remove(&self, id: &str) -> Result<()> {
660 self.db
661 .prepare("DELETE FROM repos WHERE id = ?")
662 .bind(&[id.into()])?
663 .run()
664 .await?;
665 Ok(())
666 }
667
668 /// Picks the store key for a repository about to be made, and
669 /// remembers it: the one its path gives, unless a repository already
670 /// holds that (one made in a workspace that has since been renamed,
671 /// whose old name this workspace now has), when its id.
672 ///
673 /// `namespace` is the git store namespace it goes in (shards.rs), or
674 /// `None` for the default, `default`. A name is taken in any of them.
675 pub async fn claim_store_key(&self, repo: &Repo, namespace: Option<&str>, default: &str) -> Result<String> {
676 let wanted = path_key(repo);
677 let held = self
678 .db
679 .prepare(
680 "SELECT 1 AS held FROM repos
681 WHERE store = ?1 OR (instr(store, '/') > 0 AND substr(store, instr(store, '/') + 1) = ?1)
682 UNION ALL
683 SELECT 1 AS held FROM repo_move_copies WHERE name = ?1 AND cleaned_ms IS NULL",
684 )
685 .bind(&[wanted.as_str().into()])?
686 .first::<serde_json::Value>(None)
687 .await?
688 .is_some();
689 let name = if held { repo.id.clone() } else { wanted };
690 let key = crate::shards::compose(namespace, &name, default);
691 remember_store(repo, &key);
692 Ok(key)
693 }
694
695
696 /// Moves a renamed workspace's repositories to its current slug, from
697 /// any of `stale`. A repository whose name the current slug already has
698 /// (one pushed there in the moment before this ran) stays where it is;
699 /// returns how many did.
700 pub async fn rename_namespace(&self, stale: &[String], current: &str) -> Result<usize> {
701 if stale.is_empty() {
702 return Ok(0);
703 }
704 let marks = vec!["?"; stale.len()].join(", ");
705 let mut moved: Vec<JsValue> = vec![current.into()];
706 moved.extend(stale.iter().map(|slug| JsValue::from(slug.as_str())));
707 let left: Vec<JsValue> = stale.iter().map(|slug| JsValue::from(slug.as_str())).collect();
708 let results = self
709 .db
710 .batch(vec![
711 self.db
712 .prepare(format!(
713 "UPDATE OR IGNORE repos SET namespace = ? WHERE namespace IN ({marks})"
714 ))
715 .bind(&moved)?,
716 self.db
717 .prepare(format!(
718 "SELECT count(*) AS left FROM repos WHERE namespace IN ({marks})"
719 ))
720 .bind(&left)?,
721 // Git operations follow the workspace, added together.
722 self.db
723 .prepare(format!(
724 "INSERT INTO git_operations (namespace, hour, operations)
725 SELECT ?, hour, SUM(operations) FROM git_operations WHERE namespace IN ({marks}) GROUP BY hour
726 ON CONFLICT (namespace, hour) DO UPDATE SET operations = git_operations.operations + excluded.operations"
727 ))
728 .bind(&moved)?,
729 self.db
730 .prepare(format!("DELETE FROM git_operations WHERE namespace IN ({marks})"))
731 .bind(&left)?,
732 // Paths repositories were transferred away from follow the
733 // workspace too, so the old slug's redirect then finds them.
734 self.db
735 .prepare(format!(
736 "UPDATE OR IGNORE repo_redirects SET namespace = ? WHERE namespace IN ({marks})"
737 ))
738 .bind(&moved)?,
739 ])
740 .await?;
741 #[derive(Deserialize)]
742 struct Left {
743 left: usize,
744 }
745 Ok(results
746 .get(1)
747 .map(|result| result.results::<Left>())
748 .transpose()?
749 .and_then(|rows| rows.into_iter().next())
750 .map_or(0, |row| row.left))
751 }
752
753 /// Records that the refs of the repository with this id changed, after
754 /// they did: what anything that lists them keeps goes stale.
755 pub async fn refs_moved(&self, id: &str) -> Result<()> {
756 self.bump_refs(
757 "UPDATE repos SET refs_version = refs_version + 1 WHERE id = ?
758 RETURNING refs_version, refs_open_until",
759 &[id.into()],
760 id,
761 )
762 .await
763 }
764
765 /// Records that a credential able to change the refs of the repository
766 /// with this id was handed out of g1t's hands, until `until`
767 /// (milliseconds): until then, nothing that lists them is kept.
768 pub async fn refs_open(&self, id: &str, until: u64) -> Result<()> {
769 self.bump_refs(
770 "UPDATE repos SET refs_version = refs_version + 1,
771 refs_open_until = max(coalesce(refs_open_until, 0), ?)
772 WHERE id = ? RETURNING refs_version, refs_open_until",
773 &[(until as f64).into(), id.into()],
774 id,
775 )
776 .await
777 }
778
779 async fn bump_refs(&self, sql: &str, params: &[JsValue], id: &str) -> Result<()> {
780 #[derive(Deserialize)]
781 struct Bumped {
782 refs_version: Option<f64>,
783 refs_open_until: Option<f64>,
784 }
785 let bumped = self
786 .db
787 .prepare(sql)
788 .bind(params)?
789 .first::<Bumped>(None)
790 .await?;
791 if let Some(bumped) = bumped {
792 note_refs(id, bumped.refs_version, bumped.refs_open_until);
793 }
794 Ok(())
795 }
796
797 pub async fn insert(&self, repo: &Repo) -> Result<()> {
798 self.db
799 .prepare(
800 "INSERT INTO repos
801 (id, namespace, name, description, is_private, owner_id,
802 default_branch, fork_of, created_at, store, mirror)
803 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
804 )
805 .bind(&[
806 repo.id.as_str().into(),
807 repo.namespace.as_str().into(),
808 repo.name.as_str().into(),
809 optional(&repo.description),
810 (repo.is_private as u8).into(),
811 repo.owner_id.as_str().into(),
812 repo.default_branch.as_str().into(),
813 optional(&repo.fork_of),
814 repo.created_at.as_str().into(),
815 store_key(repo).into(),
816 repo.mirror
817 .as_ref()
818 .and_then(|mirror| serde_json::to_string(mirror).ok())
819 .map_or(JsValue::NULL, JsValue::from),
820 ])?
821 .run()
822 .await?;
823 Ok(())
824 }
825}
826
827#[cfg(test)]
828mod tests {
829 use super::*;
830 use g1t_contracts::access::{BasePermission, RepoGrant};
831 use g1t_contracts::{Membership, Role, User};
832
833 #[test]
834 fn the_refs_state_kept_only_moves_forward() {
835 let mut states = RefsStates::default();
836 assert_eq!(states.get("rep_1"), None);
837 states.note("rep_1", RefsState { version: 3, open_until: 0 });
838 // A read that started before a bump and finished after it.
839 states.note("rep_1", RefsState { version: 2, open_until: 0 });
840 assert_eq!(states.get("rep_1").unwrap().version, 3);
841 states.note("rep_1", RefsState { version: 4, open_until: 9_000 });
842 states.note("rep_1", RefsState { version: 5, open_until: 0 });
843 assert_eq!(states.get("rep_1"), Some(RefsState { version: 5, open_until: 9_000 }));
844 assert_eq!(states.get("rep_2"), None);
845 }
846
847 #[test]
848 fn a_row_from_before_the_column_has_no_refs_state() {
849 let row = |version: Option<f64>| RepoRow {
850 id: format!("rep_row_{}", version.is_some()),
851 namespace: "acme".into(),
852 name: "rocket".into(),
853 description: None,
854 is_private: 0,
855 owner_id: "usr_owner".into(),
856 default_branch: "main".into(),
857 fork_of: None,
858 protected: 0,
859 created_at: String::new(),
860 store: None,
861 topics: None,
862 website: None,
863 archived_at: None,
864 mirror: None,
865 deleted_at: None,
866 refs_version: version,
867 refs_open_until: None,
868 retired_at: None,
869 retired_head: None,
870 writes_paused_until: None,
871 writes_paused_for: None,
872 };
873 let old = Repo::from(row(None));
874 assert_eq!(refs_state(&old.id), None);
875 let new = Repo::from(row(Some(7.0)));
876 assert_eq!(refs_state(&new.id), Some(RefsState { version: 7, open_until: 0 }));
877 }
878
879 #[test]
880 fn a_removed_working_copy_is_known_by_its_row() {
881 note_retired("rep_fork", Some("abc"));
882 assert_eq!(retired("rep_fork").as_deref(), Some("abc"));
883 note_retired("rep_fork", None);
884 assert_eq!(retired("rep_fork"), None);
885 }
886
887 #[test]
888 fn a_repository_read_a_moment_ago_is_reused_for_a_few_seconds() {
889 let mut recent = Recent::default();
890 let path = RepoPath {
891 namespace: "Acme".into(),
892 name: "Rocket".into(),
893 };
894 recent.keep(&path, &repo(false), 1_000);
895 // Paths are matched as the table matches them, ignoring case.
896 let lower = RepoPath {
897 namespace: "acme".into(),
898 name: "rocket".into(),
899 };
900 assert_eq!(recent.get(&lower, 1_000 + RECENT_MS - 1).unwrap().id, "rep_1");
901 assert!(recent.get(&lower, 1_000 + RECENT_MS).is_none());
902 let other = RepoPath {
903 namespace: "acme".into(),
904 name: "booster".into(),
905 };
906 assert!(recent.get(&other, 1_000).is_none());
907 // Keeping another later drops the stale row.
908 recent.keep(&other, &repo(true), 1_000 + RECENT_MS);
909 assert_eq!(recent.rows.len(), 1);
910 }
911
912 fn repo(private: bool) -> Repo {
913 Repo {
914 id: "rep_1".into(),
915 namespace: "acme".into(),
916 name: "rocket".into(),
917 description: None,
918 is_private: private,
919 owner_id: "usr_owner".into(),
920 default_branch: "main".into(),
921 fork_of: None,
922 protected: false,
923 created_at: String::new(),
924 topics: Vec::new(),
925 website: None,
926 archived_at: None,
927 mirror: None,
928 }
929 }
930
931 fn person(id: &str, memberships: Vec<Membership>, grants: Vec<(&str, RepoRole)>) -> Viewer {
932 Some(User {
933 id: id.into(),
934 username: id.into(),
935 verified: true,
936 workspaces: memberships,
937 grants: grants
938 .into_iter()
939 .map(|(repo_id, role)| RepoGrant { repo_id: repo_id.into(), workspace: "acme".into(), role, team: None })
940 .collect(),
941 ..User::default()
942 })
943 }
944
945 /// What git asks: clone and fetch need Read on a private repository,
946 /// push needs Write.
947 #[test]
948 fn git_reads_with_read_and_pushes_with_write() {
949 let private = repo(true);
950 let reader = person("usr_r", vec![], vec![("rep_1", RepoRole::Read)]);
951 assert!(can_read(&private, &reader));
952 assert!(!can_write(&private, &reader));
953 let writer = person("usr_w", vec![], vec![("rep_1", RepoRole::Write)]);
954 assert!(can_read(&private, &writer) && can_write(&private, &writer));
955 let stranger = person("usr_s", vec![], vec![("rep_2", RepoRole::Admin)]);
956 assert!(!can_read(&private, &stranger) && !can_write(&private, &stranger));
957 assert!(!can_read(&private, &None));
958 // A public repository: anyone clones, nobody without Write pushes.
959 let public = repo(false);
960 assert!(can_read(&public, &None) && !can_write(&public, &None));
961 assert!(can_read(&public, &stranger) && !can_write(&public, &stranger));
962 }
963
964 #[test]
965 fn members_follow_the_base_permission_and_owners_have_admin() {
966 let private = repo(true);
967 let default_member = person("usr_m", vec![Membership::member("acme")], vec![]);
968 assert!(can_write(&private, &default_member));
969 assert!(!can(&private, &default_member, Capability::ManageIntegrations));
970 let none = Membership { base_permission: Some(BasePermission::None), ..Membership::member("acme") };
971 let locked_out = person("usr_n", vec![none.clone()], vec![]);
972 assert!(!can_read(&private, &locked_out));
973 let given = person("usr_g", vec![none], vec![("rep_1", RepoRole::Triage)]);
974 assert!(can_read(&private, &given) && !can_write(&private, &given));
975 let owner = person("usr_o", vec![Membership { role: Role::Owner, ..Membership::member("acme") }], vec![]);
976 assert_eq!(role(&private, &owner), Some(RepoRole::Admin));
977 assert!(can(&private, &owner, Capability::Delete));
978 }
979
980 #[test]
981 fn a_pull_requests_fork_is_its_authors() {
982 let fork = Repo {
983 namespace: "pulls".into(),
984 fork_of: Some("rep_1".into()),
985 owner_id: "usr_a".into(),
986 ..repo(true)
987 };
988 let author = person("usr_a", vec![], vec![]);
989 assert!(can_write(&fork, &author));
990 assert!(!can(&fork, &author, Capability::ManageSettings));
991 let other = person("usr_b", vec![Membership::member("acme")], vec![]);
992 assert!(!can_write(&fork, &other));
993 }
994}