Skip to content

g1t/services/packages/src/db.rs

1,593 lines64,546 bytesCodeBlameRaw
1//! The packages database: packages, their versions and tags, the blobs
2//! they use, and uploads in progress.
3
4use g1t_contracts::time::rfc3339;
5use serde::Deserialize;
6use worker::wasm_bindgen::JsValue;
7use worker::{D1Database, D1PreparedStatement, Result};
8
9use crate::digest::Digest;
10use crate::upload::Progress;
11
12/// A day: how long an unreferenced blob is kept, and an upload left open.
13pub const DAY_MS: u64 = 24 * 60 * 60 * 1000;
14
15fn num(n: u64) -> JsValue {
16 JsValue::from(n as f64)
17}
18
19fn text(s: &str) -> JsValue {
20 JsValue::from(s)
21}
22
23fn opt(s: Option<&str>) -> JsValue {
24 s.map_or(JsValue::NULL, JsValue::from)
25}
26
27#[derive(Clone, Debug, Deserialize)]
28pub struct PackageRow {
29 pub id: String,
30 pub workspace: String,
31 pub ecosystem: String,
32 pub name: String,
33 pub repo_id: Option<String>,
34 pub repo_name: Option<String>,
35 pub visibility: String,
36 pub description: Option<String>,
37 pub created_by: String,
38 pub created_at: String,
39 pub updated_at: String,
40 pub downloads: u64,
41 /// Set while its workspace is deleted and may still be restored.
42 #[serde(default)]
43 pub workspace_deleted_at: Option<String>,
44 /// For a linked package: 1 when it takes its repository's roles.
45 #[serde(default = "one")]
46 pub inherit_access: u32,
47 /// Set while the package is deleted and may still be restored.
48 #[serde(default)]
49 pub deleted_at: Option<String>,
50 #[serde(default)]
51 pub deleted_by: Option<String>,
52}
53
54fn one() -> u32 {
55 1
56}
57
58impl PackageRow {
59 pub fn public(&self) -> bool {
60 self.visibility == "public"
61 }
62
63 /// Whether it is deleted, or its workspace is: hidden from everyone.
64 pub fn hidden(&self) -> bool {
65 self.workspace_deleted_at.is_some() || self.deleted_at.is_some()
66 }
67
68 pub fn inherits(&self) -> bool {
69 self.inherit_access != 0
70 }
71}
72
73/// A package with what listings add up about it.
74#[derive(Clone, Debug, Deserialize)]
75pub struct ListedRow {
76 #[serde(flatten)]
77 pub package: PackageRow,
78 pub version_count: u32,
79 pub bytes: u64,
80 pub latest_tag: Option<String>,
81 /// The version `latest_tag` points to: what an npm listing shows,
82 /// since its tags name versions rather than being what is installed.
83 #[serde(default)]
84 pub latest_tag_version: Option<String>,
85 /// Every version, newest published first, one per line: the summary
86 /// picks the highest of them (see `newest_version`).
87 pub latest_version: Option<String>,
88}
89
90/// The version a listing calls latest: the highest stable one by number
91/// (`v3.0.2` over `1.0.0`, whatever order they were published in, as an
92/// import publishes every tag at once), else the highest pre-release, else
93/// the newest published when none reads as a number.
94pub fn newest_version(versions: &str) -> Option<String> {
95 // Stable over pre-release, then by number, then pre-releases by label.
96 let parse = |version: &str| -> Option<(bool, Vec<u64>, String)> {
97 let bare = version.strip_prefix('v').unwrap_or(version);
98 let (core, pre) = match bare.split_once(['-', '+']) {
99 Some((core, rest)) if bare.as_bytes()[core.len()] == b'-' => (core, rest.to_owned()),
100 Some((core, _)) => (core, String::new()),
101 None => (bare, String::new()),
102 };
103 let parts = core.split('.').map(|part| part.parse::<u64>().ok()).collect::<Option<Vec<_>>>()?;
104 Some((pre.is_empty(), parts, pre))
105 };
106 let list: Vec<&str> = versions.lines().map(str::trim).filter(|v| !v.is_empty()).collect();
107 list.iter()
108 .filter_map(|v| parse(v).map(|key| (key, *v)))
109 .max_by(|a, b| a.0.cmp(&b.0))
110 .map(|(_, v)| v.to_owned())
111 .or_else(|| list.first().map(|v| (*v).to_owned()))
112}
113
114/// What a listing shows as a package's latest. An image's tag is what is
115/// pulled, so it is shown as is; npm's dist-tags (`latest`) name versions,
116/// so the version the tag points to is shown, as for every other registry.
117/// Without a tag, the highest version (see `newest_version`).
118pub fn latest_shown(row: &ListedRow) -> Option<String> {
119 let tagged = if row.package.ecosystem == "container" { &row.latest_tag } else { &row.latest_tag_version };
120 tagged.clone().or_else(|| row.latest_version.as_deref().and_then(newest_version))
121}
122
123#[cfg(test)]
124mod newest_tests {
125 use super::{ListedRow, PackageRow, latest_shown, newest_version};
126
127 fn listed(ecosystem: &str, tag: Option<&str>, tag_version: Option<&str>, versions: &str) -> ListedRow {
128 ListedRow {
129 package: PackageRow {
130 id: "pkg_1".into(),
131 workspace: "acme".into(),
132 ecosystem: ecosystem.into(),
133 name: "web".into(),
134 repo_id: None,
135 repo_name: None,
136 visibility: "private".into(),
137 description: None,
138 created_by: "usr_1".into(),
139 created_at: "2026-10-06T00:00:00.000Z".into(),
140 updated_at: "2026-10-06T00:00:00.000Z".into(),
141 downloads: 0,
142 workspace_deleted_at: None,
143 inherit_access: 1,
144 deleted_at: None,
145 deleted_by: None,
146 },
147 version_count: 2,
148 bytes: 0,
149 latest_tag: tag.map(str::to_owned),
150 latest_tag_version: tag_version.map(str::to_owned),
151 latest_version: Some(versions.to_owned()),
152 }
153 }
154
155 #[test]
156 fn npm_shows_the_version_its_tag_points_to_and_an_image_its_tag() {
157 let npm = listed("npm", Some("latest"), Some("1.2.0"), "2.0.0-beta.1\n1.2.0\n1.0.0");
158 assert_eq!(latest_shown(&npm).as_deref(), Some("1.2.0"), "the version, not the word latest");
159 let image = listed("container", Some("latest"), Some("sha256:abc"), "sha256:abc");
160 assert_eq!(latest_shown(&image).as_deref(), Some("latest"), "an image is pulled by its tag");
161 let cargo = listed("cargo", None, None, "0.9.0\n1.1.0\n1.0.0");
162 assert_eq!(latest_shown(&cargo).as_deref(), Some("1.1.0"), "no tags: the highest version");
163 let untagged = listed("container", None, None, "sha256:abc");
164 assert_eq!(latest_shown(&untagged).as_deref(), Some("sha256:abc"));
165 }
166
167 #[test]
168 fn the_latest_is_the_highest_stable_version_not_the_last_published() {
169 assert_eq!(newest_version("1.0.0
1703.0.2
1712.0.0
1723.0.0").as_deref(), Some("3.0.2"));
173 assert_eq!(newest_version("v1.10.0
174v1.9.3").as_deref(), Some("v1.10.0"));
175 assert_eq!(newest_version("4.0.0-beta.1
1763.0.2").as_deref(), Some("3.0.2"));
177 assert_eq!(newest_version("4.0.0-beta.1
1784.0.0-alpha").as_deref(), Some("4.0.0-beta.1"));
179 assert_eq!(newest_version("dev-main
180nightly").as_deref(), Some("dev-main"));
181 assert_eq!(newest_version(""), None);
182 }
183}
184
185#[derive(Clone, Debug, Deserialize)]
186pub struct BlobRow {
187 pub digest: String,
188 pub size: u64,
189 pub media_type: Option<String>,
190 pub object_key: String,
191}
192
193#[derive(Clone, Debug, Deserialize)]
194pub struct VersionRow {
195 pub id: String,
196 pub package_id: String,
197 pub version: String,
198 pub digest: String,
199 pub size: u64,
200 pub metadata: String,
201 pub subject: Option<String>,
202 pub published_by: Option<String>,
203 pub published_at: String,
204 /// npm's deprecation message, when the version is deprecated.
205 #[serde(default)]
206 pub deprecated: Option<String>,
207 /// Cargo: 1 when the version is yanked.
208 #[serde(default)]
209 pub yanked: u32,
210 /// Its own downloads, counted in every registry.
211 #[serde(default)]
212 pub downloads: u64,
213 /// Set while the version is deleted and may still be restored.
214 #[serde(default)]
215 pub deleted_at: Option<String>,
216 #[serde(default)]
217 pub deleted_by: Option<String>,
218}
219
220impl VersionRow {
221 pub fn is_yanked(&self) -> bool {
222 self.yanked != 0
223 }
224}
225
226impl VersionRow {
227 pub fn meta(&self) -> serde_json::Value {
228 serde_json::from_str(&self.metadata).unwrap_or_default()
229 }
230
231 pub fn media_type(&self) -> Option<String> {
232 self.meta()["media_type"].as_str().map(str::to_owned)
233 }
234}
235
236#[derive(Clone, Debug, Deserialize)]
237pub struct TagRow {
238 pub tag: String,
239 pub version_id: String,
240 pub digest: String,
241 pub updated_at: String,
242}
243
244#[derive(Clone, Debug, Deserialize)]
245pub struct UploadRow {
246 pub id: String,
247 pub workspace: String,
248 pub package_id: String,
249 pub package: String,
250 pub multipart_id: Option<String>,
251 pub parts: String,
252 pub offset: u64,
253 pub tail: u64,
254 pub hash_state: String,
255}
256
257impl UploadRow {
258 pub fn progress(&self) -> Option<Progress> {
259 Some(Progress {
260 id: self.id.clone(),
261 multipart_id: self.multipart_id.clone(),
262 parts: serde_json::from_str(&self.parts).ok()?,
263 offset: self.offset,
264 tail: self.tail,
265 hasher: crate::digest::Sha256::restore(&self.hash_state)?,
266 })
267 }
268}
269
270/// One file of a version, as it is kept.
271#[derive(Clone, Debug, Deserialize)]
272pub struct FileRow {
273 pub name: String,
274 pub digest: String,
275 pub size: u64,
276 pub media_type: Option<String>,
277}
278
279/// A file found by its name, and the package that keeps it.
280#[derive(Clone, Debug, Deserialize)]
281pub struct NamedFile {
282 pub package_id: String,
283 pub digest: String,
284 pub size: u64,
285}
286
287/// A file's other checksums, in hex, beside its SHA-256 digest.
288#[derive(Clone, Debug, PartialEq, Eq, Deserialize)]
289pub struct Checksums {
290 pub md5: String,
291 pub sha1: String,
292 pub sha512: String,
293}
294
295impl Checksums {
296 pub fn of(bytes: &[u8]) -> Checksums {
297 use sha1::Digest as _;
298 Checksums {
299 md5: format!("{:x}", md5::compute(bytes)),
300 sha1: hex::encode(sha1::Sha1::digest(bytes)),
301 sha512: hex::encode(sha2::Sha512::digest(bytes)),
302 }
303 }
304}
305
306/// One file of a version, as it is recorded.
307pub struct NewFile {
308 pub name: String,
309 pub digest: String,
310 pub size: u64,
311 pub media_type: Option<String>,
312}
313
314/// A version to record.
315pub struct NewVersion {
316 pub id: String,
317 pub package_id: String,
318 pub version: String,
319 pub digest: String,
320 pub size: u64,
321 pub metadata: String,
322 pub subject: Option<String>,
323 pub published_by: Option<String>,
324 pub files: Vec<NewFile>,
325}
326
327const PACKAGE_COLUMNS: &str = "id, workspace, ecosystem, name, repo_id, repo_name, visibility, description, created_by, created_at, updated_at, \
328 downloads, workspace_deleted_at, inherit_access, deleted_at, deleted_by";
329/// Workspaces that are deleted, waiting to be purged or restored.
330const DELETED_WORKSPACES: &str = "SELECT workspace FROM packages WHERE workspace_deleted_at IS NOT NULL";
331const VERSION_COLUMNS: &str =
332 "id, package_id, version, digest, size, metadata, subject, published_by, published_at, deprecated, yanked, downloads, deleted_at, deleted_by";
333/// A version that is not deleted: what every registry and listing sees.
334const LIVE: &str = "deleted_at IS NULL";
335
336/// A role given on a package, as kept.
337#[derive(Clone, Debug, Deserialize)]
338pub struct AccessRow {
339 pub package_id: String,
340 pub grantee_kind: String,
341 pub grantee_id: String,
342 pub grantee_name: String,
343 pub role: String,
344 pub created_at: String,
345}
346
347/// A repository given access under Manage Actions access, as kept.
348#[derive(Clone, Debug, Deserialize)]
349pub struct ActionsRow {
350 pub package_id: String,
351 pub repo_id: String,
352 pub repo_name: String,
353 pub role: String,
354 pub created_at: String,
355}
356
357pub struct Db {
358 pub db: D1Database,
359}
360
361impl Db {
362 fn prepare(&self, sql: &str, values: &[JsValue]) -> Result<D1PreparedStatement> {
363 self.db.prepare(sql).bind(values)
364 }
365
366 pub async fn package(&self, workspace: &str, ecosystem: &str, name: &str) -> Result<Option<PackageRow>> {
367 self.prepare(
368 &format!("SELECT {PACKAGE_COLUMNS} FROM packages WHERE workspace = ? AND ecosystem = ? AND name = ?"),
369 &[text(workspace), text(ecosystem), text(name)],
370 )?
371 .first(None)
372 .await
373 }
374
375 /// A package by its name in any case: Cargo's names are one name
376 /// whatever their case (`Inflector` is `inflector`).
377 pub async fn package_any_case(&self, workspace: &str, ecosystem: &str, name: &str) -> Result<Option<PackageRow>> {
378 self.prepare(
379 &format!("SELECT {PACKAGE_COLUMNS} FROM packages WHERE workspace = ? AND ecosystem = ? AND name = ? COLLATE NOCASE LIMIT 1"),
380 &[text(workspace), text(ecosystem), text(name)],
381 )?
382 .first(None)
383 .await
384 }
385
386 /// The package a new crate's name would clash with: one named the same
387 /// apart from case and `-` against `_`, as crates.io decides.
388 pub async fn package_folded(&self, workspace: &str, ecosystem: &str, folded: &str) -> Result<Option<PackageRow>> {
389 self.prepare(
390 &format!(
391 "SELECT {PACKAGE_COLUMNS} FROM packages WHERE workspace = ? AND ecosystem = ? AND replace(lower(name), '_', '-') = ? LIMIT 1"
392 ),
393 &[text(workspace), text(ecosystem), text(folded)],
394 )?
395 .first(None)
396 .await
397 }
398
399 /// Whether the workspace has a private package of the ecosystem.
400 pub async fn has_private(&self, workspace: &str, ecosystem: &str) -> Result<bool> {
401 let row: Option<serde_json::Value> = self
402 .prepare(
403 "SELECT 1 AS private FROM packages WHERE workspace = ? AND ecosystem = ? AND visibility = 'private' AND workspace_deleted_at IS NULL AND deleted_at IS NULL LIMIT 1",
404 &[text(workspace), text(ecosystem)],
405 )?
406 .first(None)
407 .await?;
408 Ok(row.is_some())
409 }
410
411 pub async fn set_yanked(&self, version_id: &str, yanked: bool) -> Result<()> {
412 self.prepare("UPDATE versions SET yanked = ? WHERE id = ?", &[num(u64::from(yanked)), text(version_id)])?
413 .run()
414 .await?;
415 Ok(())
416 }
417
418 /// Makes a package unless one of the name is already there (a push
419 /// beside this one may have made it), and answers with the one kept.
420 #[allow(clippy::too_many_arguments)]
421 pub async fn create_package(
422 &self,
423 id: &str,
424 workspace: &str,
425 ecosystem: &str,
426 name: &str,
427 repo: Option<(&str, &str, bool)>,
428 created_by: &str,
429 now_ms: u64,
430 ) -> Result<PackageRow> {
431 let now = rfc3339(now_ms);
432 let visibility = match repo {
433 Some((_, _, private)) if !private => "public",
434 _ => "private",
435 };
436 self.prepare(
437 "INSERT OR IGNORE INTO packages (id, workspace, ecosystem, name, repo_id, repo_name, visibility, created_by, created_at, updated_at)
438 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
439 &[
440 text(id),
441 text(workspace),
442 text(ecosystem),
443 text(name),
444 opt(repo.map(|r| r.0)),
445 opt(repo.map(|r| r.1)),
446 text(visibility),
447 text(created_by),
448 text(&now),
449 text(&now),
450 ],
451 )?
452 .run()
453 .await?;
454 self.package(workspace, ecosystem, name)
455 .await?
456 .ok_or_else(|| worker::Error::RustError(format!("package {workspace}/{name} was not made")))
457 }
458
459 /// A workspace's packages with their counts, newest first.
460 pub async fn list(
461 &self,
462 workspace: &str,
463 ecosystem: Option<&str>,
464 repo_id: Option<&str>,
465 query: Option<&str>,
466 limit: u32,
467 ) -> Result<Vec<ListedRow>> {
468 self.listing(workspace, ecosystem, repo_id, query, limit, false).await
469 }
470
471 /// A workspace's deleted packages, newest deletion first, with what
472 /// they held when they were deleted.
473 pub async fn deleted_packages(&self, workspace: &str, limit: u32) -> Result<Vec<ListedRow>> {
474 self.listing(workspace, None, None, None, limit, true).await
475 }
476
477 async fn listing(
478 &self,
479 workspace: &str,
480 ecosystem: Option<&str>,
481 repo_id: Option<&str>,
482 query: Option<&str>,
483 limit: u32,
484 deleted: bool,
485 ) -> Result<Vec<ListedRow>> {
486 let mut sql = format!(
487 "SELECT {}, \
488 (SELECT COUNT(*) FROM versions v WHERE v.package_id = p.id AND v.deleted_at IS NULL) AS version_count, \
489 (SELECT COALESCE(SUM(b.size), 0) FROM blobs b WHERE b.digest IN \
490 (SELECT vf.digest FROM version_files vf JOIN versions v ON v.id = vf.version_id \
491 WHERE v.package_id = p.id AND v.deleted_at IS NULL)) AS bytes, \
492 (SELECT t.tag FROM tags t JOIN versions v ON v.id = t.version_id WHERE t.package_id = p.id AND v.deleted_at IS NULL \
493 ORDER BY t.tag = 'latest' DESC, t.updated_at DESC LIMIT 1) AS latest_tag, \
494 (SELECT v.version FROM tags t JOIN versions v ON v.id = t.version_id WHERE t.package_id = p.id AND v.deleted_at IS NULL \
495 ORDER BY t.tag = 'latest' DESC, t.updated_at DESC LIMIT 1) AS latest_tag_version, \
496 (SELECT GROUP_CONCAT(version, char(10)) FROM \
497 (SELECT v.version FROM versions v WHERE v.package_id = p.id AND v.yanked = 0 AND v.deleted_at IS NULL \
498 ORDER BY v.published_at DESC)) AS latest_version \
499 FROM packages p WHERE p.workspace = ? AND p.workspace_deleted_at IS NULL AND p.deleted_at IS {}",
500 PACKAGE_COLUMNS.split(", ").map(|c| format!("p.{c}")).collect::<Vec<_>>().join(", "),
501 if deleted { "NOT NULL" } else { "NULL" }
502 );
503 let mut values = vec![text(workspace)];
504 if let Some(ecosystem) = ecosystem {
505 sql.push_str(" AND p.ecosystem = ?");
506 values.push(text(ecosystem));
507 }
508 if let Some(repo_id) = repo_id {
509 sql.push_str(" AND p.repo_id = ?");
510 values.push(text(repo_id));
511 }
512 if let Some(query) = query.map(str::trim).filter(|q| !q.is_empty()) {
513 sql.push_str(" AND p.name LIKE ? ESCAPE '\\'");
514 let escaped = query.replace('\\', "\\\\").replace('%', "\\%").replace('_', "\\_");
515 values.push(text(&format!("%{}%", escaped.to_lowercase())));
516 }
517 let order = if deleted { "p.deleted_at" } else { "p.updated_at" };
518 sql.push_str(&format!(" ORDER BY {order} DESC LIMIT {limit}"));
519 self.prepare(&sql, &values)?.all().await?.results()
520 }
521
522 pub async fn set_visibility(&self, package_id: &str, visibility: &str, now_ms: u64) -> Result<()> {
523 self.prepare(
524 "UPDATE packages SET visibility = ?, updated_at = ? WHERE id = ?",
525 &[text(visibility), text(&rfc3339(now_ms)), text(package_id)],
526 )?
527 .run()
528 .await?;
529 Ok(())
530 }
531
532 pub async fn set_link(&self, package_id: &str, repo: Option<(&str, &str)>, visibility: &str, now_ms: u64) -> Result<()> {
533 self.prepare(
534 "UPDATE packages SET repo_id = ?, repo_name = ?, visibility = ?, updated_at = ? WHERE id = ?",
535 &[opt(repo.map(|r| r.0)), opt(repo.map(|r| r.1)), text(visibility), text(&rfc3339(now_ms)), text(package_id)],
536 )?
537 .run()
538 .await?;
539 Ok(())
540 }
541
542 /// Adds downloads to packages, and to the versions named with them.
543 pub async fn add_downloads(&self, counts: &[((String, Option<String>), u64)]) -> Result<()> {
544 if counts.is_empty() {
545 return Ok(());
546 }
547 let mut batch = Vec::with_capacity(counts.len());
548 for ((id, version), count) in counts {
549 batch.push(self.prepare("UPDATE packages SET downloads = downloads + ? WHERE id = ?", &[num(*count), text(id)])?);
550 if let Some(version) = version {
551 batch.push(self.prepare("UPDATE versions SET downloads = downloads + ? WHERE id = ?", &[num(*count), text(version)])?);
552 }
553 }
554 self.db.batch(batch).await?;
555 Ok(())
556 }
557
558 /// A package and everything it holds, for good: the purge, and a
559 /// Composer package whose repository went. Its blobs are left to the
560 /// sweep.
561 pub async fn delete_package(&self, package_id: &str) -> Result<()> {
562 let id = [text(package_id)];
563 self.db
564 .batch(vec![
565 self.prepare("DELETE FROM tags WHERE package_id = ?", &id)?,
566 self.prepare("DELETE FROM version_files WHERE version_id IN (SELECT id FROM versions WHERE package_id = ?)", &id)?,
567 self.prepare("DELETE FROM versions WHERE package_id = ?", &id)?,
568 self.prepare("DELETE FROM package_blobs WHERE package_id = ?", &id)?,
569 self.prepare("DELETE FROM uploads WHERE package_id = ?", &id)?,
570 self.prepare("DELETE FROM package_access WHERE package_id = ?", &id)?,
571 self.prepare("DELETE FROM package_actions_access WHERE package_id = ?", &id)?,
572 self.prepare("DELETE FROM packages WHERE id = ?", &id)?,
573 ])
574 .await?;
575 Ok(())
576 }
577
578 pub async fn blob(&self, digest: &Digest) -> Result<Option<BlobRow>> {
579 self.prepare("SELECT digest, size, media_type, object_key FROM blobs WHERE digest = ?", &[text(digest.as_str())])?
580 .first(None)
581 .await
582 }
583
584 /// The blob, if this package may serve it.
585 pub async fn package_blob(&self, package_id: &str, digest: &Digest) -> Result<Option<BlobRow>> {
586 self.prepare(
587 "SELECT b.digest, b.size, b.media_type, b.object_key FROM package_blobs pb JOIN blobs b ON b.digest = pb.digest
588 WHERE pb.package_id = ? AND pb.digest = ?",
589 &[text(package_id), text(digest.as_str())],
590 )?
591 .first(None)
592 .await
593 }
594
595 /// Records a stored blob, or where it is now, and lets the package
596 /// serve it. Either way the blob is marked as just used.
597 pub async fn keep_blob(&self, package_id: &str, digest: &Digest, size: u64, media_type: Option<&str>, key: &str, now_ms: u64) -> Result<()> {
598 let now = rfc3339(now_ms);
599 self.db
600 .batch(vec![
601 self.prepare(
602 "INSERT INTO blobs (digest, size, media_type, object_key, created_at, touched_at) VALUES (?, ?, ?, ?, ?, ?)
603 ON CONFLICT (digest) DO UPDATE SET object_key = excluded.object_key, size = excluded.size",
604 &[text(digest.as_str()), num(size), opt(media_type), text(key), text(&now), text(&now)],
605 )?,
606 self.prepare("UPDATE blobs SET touched_at = ? WHERE digest = ?", &[text(&now), text(digest.as_str())])?,
607 self.prepare(
608 "INSERT OR IGNORE INTO package_blobs (package_id, digest, created_at) VALUES (?, ?, ?)",
609 &[text(package_id), text(digest.as_str()), text(&now)],
610 )?,
611 ])
612 .await?;
613 Ok(())
614 }
615
616 /// Lets the package serve a blob that is already stored.
617 pub async fn link_blob(&self, package_id: &str, digest: &Digest, now_ms: u64) -> Result<()> {
618 let now = rfc3339(now_ms);
619 self.db
620 .batch(vec![
621 self.prepare("UPDATE blobs SET touched_at = ? WHERE digest = ?", &[text(&now), text(digest.as_str())])?,
622 self.prepare(
623 "INSERT OR IGNORE INTO package_blobs (package_id, digest, created_at) VALUES (?, ?, ?)",
624 &[text(package_id), text(digest.as_str()), text(&now)],
625 )?,
626 ])
627 .await?;
628 Ok(())
629 }
630
631 /// Whether any version of the package names the blob.
632 pub async fn blob_in_use(&self, package_id: &str, digest: &Digest) -> Result<bool> {
633 let row: Option<serde_json::Value> = self
634 .prepare(
635 "SELECT 1 AS used FROM version_files vf JOIN versions v ON v.id = vf.version_id WHERE v.package_id = ? AND vf.digest = ? LIMIT 1",
636 &[text(package_id), text(digest.as_str())],
637 )?
638 .first(None)
639 .await?;
640 Ok(row.is_some())
641 }
642
643 pub async fn unlink_blob(&self, package_id: &str, digest: &Digest) -> Result<()> {
644 self.prepare("DELETE FROM package_blobs WHERE package_id = ? AND digest = ?", &[text(package_id), text(digest.as_str())])?
645 .run()
646 .await?;
647 Ok(())
648 }
649
650 pub async fn create_upload(&self, row: &UploadRow, now_ms: u64) -> Result<()> {
651 self.prepare(
652 "INSERT INTO uploads (id, workspace, package_id, package, parts, \"offset\", tail, hash_state, created_at, expires_at)
653 VALUES (?, ?, ?, ?, '[]', 0, 0, ?, ?, ?)",
654 &[
655 text(&row.id),
656 text(&row.workspace),
657 text(&row.package_id),
658 text(&row.package),
659 text(&row.hash_state),
660 text(&rfc3339(now_ms)),
661 text(&rfc3339(now_ms + DAY_MS)),
662 ],
663 )?
664 .run()
665 .await?;
666 Ok(())
667 }
668
669 pub async fn upload(&self, id: &str) -> Result<Option<UploadRow>> {
670 self.prepare(
671 "SELECT id, workspace, package_id, package, multipart_id, parts, \"offset\" AS offset, tail, hash_state FROM uploads WHERE id = ?",
672 &[text(id)],
673 )?
674 .first(None)
675 .await
676 }
677
678 pub async fn save_progress(&self, progress: &Progress, now_ms: u64) -> Result<()> {
679 self.prepare(
680 "UPDATE uploads SET multipart_id = ?, parts = ?, \"offset\" = ?, tail = ?, hash_state = ?, expires_at = ? WHERE id = ?",
681 &[
682 opt(progress.multipart_id.as_deref()),
683 text(&serde_json::to_string(&progress.parts)?),
684 num(progress.offset),
685 num(progress.tail),
686 text(&progress.hasher.save()),
687 text(&rfc3339(now_ms + DAY_MS)),
688 text(&progress.id),
689 ],
690 )?
691 .run()
692 .await?;
693 Ok(())
694 }
695
696 pub async fn delete_upload(&self, id: &str) -> Result<()> {
697 self.prepare("DELETE FROM uploads WHERE id = ?", &[text(id)])?.run().await?;
698 Ok(())
699 }
700
701 pub async fn expired_uploads(&self, now_ms: u64, limit: u32) -> Result<Vec<UploadRow>> {
702 self.prepare(
703 &format!(
704 "SELECT id, workspace, package_id, package, multipart_id, parts, \"offset\" AS offset, tail, hash_state
705 FROM uploads WHERE expires_at < ? ORDER BY expires_at LIMIT {limit}"
706 ),
707 &[text(&rfc3339(now_ms))],
708 )?
709 .all()
710 .await?
711 .results()
712 }
713
714 pub async fn version_by_digest(&self, package_id: &str, digest: &str) -> Result<Option<VersionRow>> {
715 self.prepare(
716 &format!("SELECT {VERSION_COLUMNS} FROM versions WHERE package_id = ? AND digest = ? AND {LIVE} LIMIT 1"),
717 &[text(package_id), text(digest)],
718 )?
719 .first(None)
720 .await
721 }
722
723 pub async fn version_by_tag(&self, package_id: &str, tag: &str) -> Result<Option<VersionRow>> {
724 self.prepare(
725 &format!(
726 "SELECT {} FROM tags t JOIN versions v ON v.id = t.version_id WHERE t.package_id = ? AND t.tag = ? AND v.deleted_at IS NULL",
727 VERSION_COLUMNS.split(", ").map(|c| format!("v.{c}")).collect::<Vec<_>>().join(", ")
728 ),
729 &[text(package_id), text(tag)],
730 )?
731 .first(None)
732 .await
733 }
734
735 /// A version by its version, digest, or a tag that points to it.
736 pub async fn find_version(&self, package_id: &str, reference: &str) -> Result<Option<VersionRow>> {
737 if reference.starts_with("ver_") {
738 let by_id: Option<VersionRow> = self
739 .prepare(
740 &format!("SELECT {VERSION_COLUMNS} FROM versions WHERE package_id = ? AND id = ? AND {LIVE}"),
741 &[text(package_id), text(reference)],
742 )?
743 .first(None)
744 .await?;
745 if by_id.is_some() {
746 return Ok(by_id);
747 }
748 }
749 if let Some(found) = self.version_by_digest(package_id, reference).await? {
750 return Ok(Some(found));
751 }
752 let by_version: Option<VersionRow> = self
753 .prepare(
754 &format!("SELECT {VERSION_COLUMNS} FROM versions WHERE package_id = ? AND version = ? AND {LIVE}"),
755 &[text(package_id), text(reference)],
756 )?
757 .first(None)
758 .await?;
759 if by_version.is_some() {
760 return Ok(by_version);
761 }
762 self.version_by_tag(package_id, reference).await
763 }
764
765 pub async fn versions(&self, package_id: &str, limit: u32) -> Result<Vec<VersionRow>> {
766 self.prepare(
767 &format!("SELECT {VERSION_COLUMNS} FROM versions WHERE package_id = ? AND {LIVE} ORDER BY published_at DESC, id DESC LIMIT {limit}"),
768 &[text(package_id)],
769 )?
770 .all()
771 .await?
772 .results()
773 }
774
775 /// Records a version unless one of its digest is there, moves `tag` to
776 /// it, and says whether anything was published: a new version, or a
777 /// tag that now points somewhere else.
778 pub async fn publish(&self, version: NewVersion, tag: Option<&str>, now_ms: u64) -> Result<(VersionRow, bool)> {
779 let now = rfc3339(now_ms);
780 let existing = self.version_by_digest(&version.package_id, &version.digest).await?;
781 let mut changed = existing.is_none();
782 let mut batch = Vec::new();
783 let version_id = match &existing {
784 Some(row) => row.id.clone(),
785 None => {
786 batch.push(self.prepare(
787 "INSERT INTO versions (id, package_id, version, digest, size, metadata, subject, published_by, published_at)
788 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)",
789 &[
790 text(&version.id),
791 text(&version.package_id),
792 text(&version.version),
793 text(&version.digest),
794 num(version.size),
795 text(&version.metadata),
796 opt(version.subject.as_deref()),
797 opt(version.published_by.as_deref()),
798 text(&now),
799 ],
800 )?);
801 for file in &version.files {
802 batch.push(self.prepare(
803 "INSERT OR IGNORE INTO version_files (version_id, name, digest, size, media_type) VALUES (?, ?, ?, ?, ?)",
804 &[text(&version.id), text(&file.name), text(&file.digest), num(file.size), opt(file.media_type.as_deref())],
805 )?);
806 }
807 version.id.clone()
808 }
809 };
810 if let Some(tag) = tag {
811 let before = self.version_by_tag(&version.package_id, tag).await?;
812 changed |= before.map(|row| row.id) != Some(version_id.clone());
813 batch.push(self.prepare(
814 "INSERT INTO tags (package_id, tag, version_id, updated_at) VALUES (?, ?, ?, ?)
815 ON CONFLICT (package_id, tag) DO UPDATE SET version_id = excluded.version_id, updated_at = excluded.updated_at",
816 &[text(&version.package_id), text(tag), text(&version_id), text(&now)],
817 )?);
818 }
819 if changed {
820 batch.push(self.prepare("UPDATE packages SET updated_at = ? WHERE id = ?", &[text(&now), text(&version.package_id)])?);
821 }
822 if !batch.is_empty() {
823 self.db.batch(batch).await?;
824 }
825 let row = self
826 .version_by_digest(&version.package_id, &version.digest)
827 .await?
828 .ok_or_else(|| worker::Error::RustError("the version was not recorded".into()))?;
829 Ok((row, changed))
830 }
831
832 pub async fn tags(&self, package_id: &str) -> Result<Vec<TagRow>> {
833 self.prepare(
834 "SELECT t.tag, t.version_id, v.digest, t.updated_at FROM tags t JOIN versions v ON v.id = t.version_id
835 WHERE t.package_id = ? AND v.deleted_at IS NULL ORDER BY t.tag",
836 &[text(package_id)],
837 )?
838 .all()
839 .await?
840 .results()
841 }
842
843 /// Tag names in order, `n` of them after `last`.
844 pub async fn tag_names(&self, package_id: &str, last: Option<&str>, n: u32) -> Result<Vec<String>> {
845 #[derive(Deserialize)]
846 struct Name {
847 tag: String,
848 }
849 let rows: Vec<Name> = self
850 .prepare(
851 &format!("SELECT t.tag FROM tags t JOIN versions v ON v.id = t.version_id WHERE t.package_id = ? AND t.tag > ? AND v.deleted_at IS NULL ORDER BY t.tag LIMIT {n}"),
852 &[text(package_id), text(last.unwrap_or(""))],
853 )?
854 .all()
855 .await?
856 .results()?;
857 Ok(rows.into_iter().map(|row| row.tag).collect())
858 }
859
860 pub async fn referrers(&self, package_id: &str, subject: &str) -> Result<Vec<VersionRow>> {
861 self.prepare(
862 &format!("SELECT {VERSION_COLUMNS} FROM versions WHERE package_id = ? AND subject = ? AND {LIVE} ORDER BY published_at"),
863 &[text(package_id), text(subject)],
864 )?
865 .all()
866 .await?
867 .results()
868 }
869
870 /// The package of an ecosystem built from a repository (Composer's).
871 pub async fn package_for_repo(&self, repo_id: &str, ecosystem: &str) -> Result<Option<PackageRow>> {
872 self.prepare(
873 &format!("SELECT {PACKAGE_COLUMNS} FROM packages WHERE repo_id = ? AND ecosystem = ? LIMIT 1"),
874 &[text(repo_id), text(ecosystem)],
875 )?
876 .first(None)
877 .await
878 }
879
880 pub async fn rename_package(&self, package_id: &str, name: &str, now_ms: u64) -> Result<()> {
881 self.prepare(
882 "UPDATE packages SET name = ?, updated_at = ? WHERE id = ?",
883 &[text(name), text(&rfc3339(now_ms)), text(package_id)],
884 )?
885 .run()
886 .await?;
887 Ok(())
888 }
889
890 /// Records a version, in place of one of the same version string
891 /// (a tag or branch that moved): its files and tags go with the old one.
892 pub async fn replace_version(&self, version: &NewVersion, now_ms: u64) -> Result<()> {
893 let now = rfc3339(now_ms);
894 let old = [text(&version.package_id), text(&version.version)];
895 let old_ids = "SELECT id FROM versions WHERE package_id = ? AND version = ?";
896 let mut batch = vec![
897 self.prepare(&format!("DELETE FROM tags WHERE version_id IN ({old_ids})"), &old)?,
898 self.prepare(&format!("DELETE FROM version_files WHERE version_id IN ({old_ids})"), &old)?,
899 self.prepare("DELETE FROM versions WHERE package_id = ? AND version = ?", &old)?,
900 self.prepare(
901 "INSERT INTO versions (id, package_id, version, digest, size, metadata, subject, published_by, published_at)
902 VALUES (?, ?, ?, ?, ?, ?, NULL, ?, ?)",
903 &[
904 text(&version.id),
905 text(&version.package_id),
906 text(&version.version),
907 text(&version.digest),
908 num(version.size),
909 text(&version.metadata),
910 opt(version.published_by.as_deref()),
911 text(&now),
912 ],
913 )?,
914 ];
915 for file in &version.files {
916 batch.push(self.prepare(
917 "INSERT OR IGNORE INTO version_files (version_id, name, digest, size, media_type) VALUES (?, ?, ?, ?, ?)",
918 &[text(&version.id), text(&file.name), text(&file.digest), num(file.size), opt(file.media_type.as_deref())],
919 )?);
920 }
921 batch.push(self.prepare("UPDATE packages SET updated_at = ? WHERE id = ?", &[text(&now), text(&version.package_id)])?);
922 self.db.batch(batch).await?;
923 Ok(())
924 }
925
926 /// The zip already made of a commit of the package, if one was.
927 pub async fn dist_for_commit(&self, package_id: &str, commit: &str) -> Result<Option<BlobRow>> {
928 self.prepare(
929 "SELECT b.digest, b.size, b.media_type, b.object_key FROM versions v
930 JOIN version_files vf ON vf.version_id = v.id AND vf.name = 'dist'
931 JOIN blobs b ON b.digest = vf.digest
932 WHERE v.package_id = ? AND v.digest = ? LIMIT 1",
933 &[text(package_id), text(commit)],
934 )?
935 .first(None)
936 .await
937 }
938
939 /// Records the zip of a commit as a file of every version at it.
940 pub async fn add_dist(&self, package_id: &str, commit: &str, digest: &Digest, size: u64) -> Result<()> {
941 let at = [text(package_id), text(commit)];
942 self.db
943 .batch(vec![
944 self.prepare(
945 "INSERT OR IGNORE INTO version_files (version_id, name, digest, size, media_type)
946 SELECT id, 'dist', ?, ?, 'application/zip' FROM versions WHERE package_id = ? AND digest = ?",
947 &[text(digest.as_str()), num(size), text(package_id), text(commit)],
948 )?,
949 self.prepare("UPDATE versions SET size = ? WHERE package_id = ? AND digest = ?", &[num(size), at[0].clone(), at[1].clone()])?,
950 ])
951 .await?;
952 Ok(())
953 }
954
955 /// Where the Composer backfill is: the last repository done, and
956 /// whether it went through them all.
957 pub async fn backfill(&self) -> Result<(Option<String>, bool)> {
958 #[derive(Deserialize)]
959 struct Row {
960 after: Option<String>,
961 finished_at: Option<String>,
962 }
963 let row: Option<Row> = self
964 .prepare("SELECT after, finished_at FROM composer_backfill WHERE key = 'repos'", &[])?
965 .first(None)
966 .await?;
967 Ok(row.map_or((None, false), |row| (row.after, row.finished_at.is_some())))
968 }
969
970 pub async fn set_backfill(&self, after: Option<&str>, finished: bool, now_ms: u64) -> Result<()> {
971 self.prepare(
972 "INSERT INTO composer_backfill (key, after, finished_at) VALUES ('repos', ?, ?)
973 ON CONFLICT (key) DO UPDATE SET after = excluded.after, finished_at = excluded.finished_at",
974 &[opt(after), if finished { text(&rfc3339(now_ms)) } else { JsValue::NULL }],
975 )?
976 .run()
977 .await?;
978 Ok(())
979 }
980
981 /// Points `tag` at a version, made or moved.
982 pub async fn set_tag(&self, package_id: &str, tag: &str, version_id: &str, now_ms: u64) -> Result<()> {
983 self.prepare(
984 "INSERT INTO tags (package_id, tag, version_id, updated_at) VALUES (?, ?, ?, ?)
985 ON CONFLICT (package_id, tag) DO UPDATE SET version_id = excluded.version_id, updated_at = excluded.updated_at",
986 &[text(package_id), text(tag), text(version_id), text(&rfc3339(now_ms))],
987 )?
988 .run()
989 .await?;
990 Ok(())
991 }
992
993 /// A version by its version string alone.
994 pub async fn version_named(&self, package_id: &str, version: &str) -> Result<Option<VersionRow>> {
995 self.prepare(
996 &format!("SELECT {VERSION_COLUMNS} FROM versions WHERE package_id = ? AND version = ? AND {LIVE}"),
997 &[text(package_id), text(version)],
998 )?
999 .first(None)
1000 .await
1001 }
1002
1003 pub async fn set_deprecated(&self, version_id: &str, message: Option<&str>) -> Result<()> {
1004 self.prepare("UPDATE versions SET deprecated = ? WHERE id = ?", &[opt(message), text(version_id)])?
1005 .run()
1006 .await?;
1007 Ok(())
1008 }
1009
1010 /// The README shown for the package, kept as a blob, and its description.
1011 pub async fn set_readme(&self, package_id: &str, digest: Option<&str>, description: Option<&str>, now_ms: u64) -> Result<()> {
1012 self.prepare(
1013 "UPDATE packages SET readme_digest = ?, description = ?, updated_at = ? WHERE id = ?",
1014 &[opt(digest), opt(description), text(&rfc3339(now_ms)), text(package_id)],
1015 )?
1016 .run()
1017 .await?;
1018 Ok(())
1019 }
1020
1021 pub async fn readme_digest(&self, package_id: &str) -> Result<Option<String>> {
1022 #[derive(Deserialize)]
1023 struct Readme {
1024 readme_digest: Option<String>,
1025 }
1026 let row: Option<Readme> = self
1027 .prepare("SELECT readme_digest FROM packages WHERE id = ?", &[text(package_id)])?
1028 .first(None)
1029 .await?;
1030 Ok(row.and_then(|r| r.readme_digest))
1031 }
1032
1033 pub async fn touch_package(&self, package_id: &str, now_ms: u64) -> Result<()> {
1034 self.prepare("UPDATE packages SET updated_at = ? WHERE id = ?", &[text(&rfc3339(now_ms)), text(package_id)])?
1035 .run()
1036 .await?;
1037 Ok(())
1038 }
1039
1040 pub async fn delete_tag(&self, package_id: &str, tag: &str) -> Result<()> {
1041 self.prepare("DELETE FROM tags WHERE package_id = ? AND tag = ?", &[text(package_id), text(tag)])?
1042 .run()
1043 .await?;
1044 Ok(())
1045 }
1046
1047 /// A version, its tags and its files, for good: the purge, and Composer's
1048 /// versions, which follow git. Blobs are left to the sweep.
1049 pub async fn delete_version(&self, version_id: &str) -> Result<()> {
1050 let id = [text(version_id)];
1051 self.db
1052 .batch(vec![
1053 self.prepare("DELETE FROM tags WHERE version_id = ?", &id)?,
1054 self.prepare("DELETE FROM version_files WHERE version_id = ?", &id)?,
1055 self.prepare("DELETE FROM versions WHERE id = ?", &id)?,
1056 ])
1057 .await?;
1058 Ok(())
1059 }
1060
1061 /// Works out again what the workspace stores: each blob its versions
1062 /// use, once, public when any public package uses it.
1063 pub async fn measure(&self, workspace: &str) -> Result<()> {
1064 let ws = [text(workspace)];
1065 self.db
1066 .batch(vec![
1067 self.prepare("DELETE FROM workspace_blobs WHERE workspace = ?", &ws)?,
1068 self.prepare(
1069 "INSERT INTO workspace_blobs (workspace, digest, size, public)
1070 SELECT p.workspace, vf.digest, MAX(vf.size), MAX(p.visibility = 'public')
1071 FROM version_files vf JOIN versions v ON v.id = vf.version_id JOIN packages p ON p.id = v.package_id
1072 WHERE p.workspace = ? AND p.deleted_at IS NULL AND v.deleted_at IS NULL GROUP BY vf.digest",
1073 &ws,
1074 )?,
1075 ])
1076 .await?;
1077 Ok(())
1078 }
1079
1080 pub async fn storage(&self, workspace: &str) -> Result<(u64, u64)> {
1081 #[derive(Deserialize)]
1082 struct Sums {
1083 public_bytes: Option<u64>,
1084 private_bytes: Option<u64>,
1085 }
1086 let sums: Option<Sums> = self
1087 .prepare(
1088 &format!(
1089 "SELECT SUM(CASE WHEN public = 1 THEN size ELSE 0 END) AS public_bytes,
1090 SUM(CASE WHEN public = 1 THEN 0 ELSE size END) AS private_bytes
1091 FROM workspace_blobs WHERE workspace = ? AND workspace NOT IN ({DELETED_WORKSPACES})"
1092 ),
1093 &[text(workspace)],
1094 )?
1095 .first(None)
1096 .await?;
1097 Ok(sums.map_or((0, 0), |s| (s.public_bytes.unwrap_or(0), s.private_bytes.unwrap_or(0))))
1098 }
1099
1100 /// Hides (`deleted_at` given) or shows again (`None`) a workspace's
1101 /// packages. Hiding keeps a mark already set; showing clears only marks.
1102 pub async fn mark_workspace(&self, workspace: &str, deleted_at: Option<&str>) -> Result<()> {
1103 let statement = match deleted_at {
1104 Some(at) => self.prepare(
1105 "UPDATE packages SET workspace_deleted_at = ? WHERE workspace = ? AND workspace_deleted_at IS NULL",
1106 &[text(at), text(workspace)],
1107 )?,
1108 None => self.prepare(
1109 "UPDATE packages SET workspace_deleted_at = NULL WHERE workspace = ? AND workspace_deleted_at IS NOT NULL",
1110 &[text(workspace)],
1111 )?,
1112 };
1113 statement.run().await?;
1114 Ok(())
1115 }
1116
1117 /// Whether the workspace is deleted, as its packages say.
1118 pub async fn workspace_hidden(&self, workspace: &str) -> Result<bool> {
1119 let row: Option<serde_json::Value> = self
1120 .prepare("SELECT 1 AS hidden FROM packages WHERE workspace = ? AND workspace_deleted_at IS NOT NULL LIMIT 1", &[text(workspace)])?
1121 .first(None)
1122 .await?;
1123 Ok(row.is_some())
1124 }
1125
1126 /// Every workspace's package storage, by workspace.
1127 pub async fn storage_all(&self) -> Result<Vec<g1t_contracts::packages::WorkspacePackageStorage>> {
1128 self.prepare(
1129 &format!(
1130 "SELECT workspace,
1131 COALESCE(SUM(CASE WHEN public = 1 THEN size ELSE 0 END), 0) AS public_bytes,
1132 COALESCE(SUM(CASE WHEN public = 1 THEN 0 ELSE size END), 0) AS private_bytes
1133 FROM workspace_blobs WHERE workspace NOT IN ({DELETED_WORKSPACES}) GROUP BY workspace ORDER BY workspace"
1134 ),
1135 &[],
1136 )?
1137 .all()
1138 .await?
1139 .results()
1140 }
1141
1142 /// Which of `digests` the workspace already holds.
1143 pub async fn held(&self, workspace: &str, digests: &[String]) -> Result<std::collections::HashSet<String>> {
1144 #[derive(Deserialize)]
1145 struct Held {
1146 digest: String,
1147 }
1148 let mut held = std::collections::HashSet::new();
1149 for chunk in digests.chunks(50) {
1150 let marks = vec!["?"; chunk.len()].join(", ");
1151 let mut values = vec![text(workspace)];
1152 values.extend(chunk.iter().map(|d| text(d)));
1153 let rows: Vec<Held> = self
1154 .prepare(&format!("SELECT digest FROM workspace_blobs WHERE workspace = ? AND digest IN ({marks})"), &values)?
1155 .all()
1156 .await?
1157 .results()?;
1158 held.extend(rows.into_iter().map(|row| row.digest));
1159 }
1160 Ok(held)
1161 }
1162
1163 /// Blobs no version uses and nothing has touched for a day.
1164 pub async fn unused_blobs(&self, now_ms: u64, limit: u32) -> Result<Vec<BlobRow>> {
1165 self.prepare(
1166 &format!(
1167 "SELECT digest, size, media_type, object_key FROM blobs b
1168 WHERE touched_at < ? AND NOT EXISTS (SELECT 1 FROM version_files vf WHERE vf.digest = b.digest)
1169 AND NOT EXISTS (SELECT 1 FROM packages p WHERE p.readme_digest = b.digest)
1170 ORDER BY touched_at LIMIT {limit}"
1171 ),
1172 &[text(&rfc3339(now_ms.saturating_sub(DAY_MS)))],
1173 )?
1174 .all()
1175 .await?
1176 .results()
1177 }
1178
1179 /// A version by its version string, made from `version` when there is
1180 /// none, and whether it was made now. Files are added one at a time
1181 /// with `put_file` (Maven uploads each in a request of its own).
1182 pub async fn version_or_new(&self, version: &NewVersion, now_ms: u64) -> Result<(VersionRow, bool)> {
1183 self.prepare(
1184 "INSERT OR IGNORE INTO versions (id, package_id, version, digest, size, metadata, subject, published_by, published_at)
1185 VALUES (?, ?, ?, ?, 0, ?, NULL, ?, ?)",
1186 &[
1187 text(&version.id),
1188 text(&version.package_id),
1189 text(&version.version),
1190 text(&version.digest),
1191 text(&version.metadata),
1192 opt(version.published_by.as_deref()),
1193 text(&rfc3339(now_ms)),
1194 ],
1195 )?
1196 .run()
1197 .await?;
1198 let row = self
1199 .version_named(&version.package_id, &version.version)
1200 .await?
1201 .ok_or_else(|| worker::Error::RustError("the version was not recorded".into()))?;
1202 let made = row.id == version.id;
1203 Ok((row, made))
1204 }
1205
1206 /// Adds a file to a version, or replaces the one of its name, and
1207 /// works out the version's size again.
1208 pub async fn put_file(&self, package_id: &str, version_id: &str, file: &NewFile, now_ms: u64) -> Result<()> {
1209 let now = rfc3339(now_ms);
1210 self.db
1211 .batch(vec![
1212 self.prepare(
1213 "INSERT INTO version_files (version_id, name, digest, size, media_type) VALUES (?, ?, ?, ?, ?)
1214 ON CONFLICT (version_id, name) DO UPDATE SET digest = excluded.digest, size = excluded.size, media_type = excluded.media_type",
1215 &[text(version_id), text(&file.name), text(&file.digest), num(file.size), opt(file.media_type.as_deref())],
1216 )?,
1217 self.prepare(
1218 "UPDATE versions SET size = (SELECT COALESCE(SUM(size), 0) FROM version_files WHERE version_id = ?) WHERE id = ?",
1219 &[text(version_id), text(version_id)],
1220 )?,
1221 self.prepare("UPDATE packages SET updated_at = ? WHERE id = ?", &[text(&now), text(package_id)])?,
1222 ])
1223 .await?;
1224 Ok(())
1225 }
1226
1227 /// A version's files, by name.
1228 pub async fn files(&self, version_id: &str) -> Result<Vec<FileRow>> {
1229 self.prepare("SELECT name, digest, size, media_type FROM version_files WHERE version_id = ? ORDER BY name", &[text(version_id)])?
1230 .all()
1231 .await?
1232 .results()
1233 }
1234
1235 pub async fn file(&self, version_id: &str, name: &str) -> Result<Option<FileRow>> {
1236 self.prepare(
1237 "SELECT name, digest, size, media_type FROM version_files WHERE version_id = ? AND name = ?",
1238 &[text(version_id), text(name)],
1239 )?
1240 .first(None)
1241 .await
1242 }
1243
1244 /// Changes what a version is named by and keeps about itself.
1245 pub async fn set_version(&self, version_id: &str, digest: &str, metadata: &str) -> Result<()> {
1246 self.prepare("UPDATE versions SET digest = ?, metadata = ? WHERE id = ?", &[text(digest), text(metadata), text(version_id)])?
1247 .run()
1248 .await?;
1249 Ok(())
1250 }
1251
1252 pub async fn set_checksums(&self, digest: &Digest, sums: &Checksums) -> Result<()> {
1253 self.prepare(
1254 "INSERT OR IGNORE INTO checksums (digest, md5, sha1, sha512) VALUES (?, ?, ?, ?)",
1255 &[text(digest.as_str()), text(&sums.md5), text(&sums.sha1), text(&sums.sha512)],
1256 )?
1257 .run()
1258 .await?;
1259 Ok(())
1260 }
1261
1262 pub async fn checksums(&self, digest: &Digest) -> Result<Option<Checksums>> {
1263 self.prepare("SELECT md5, sha1, sha512 FROM checksums WHERE digest = ?", &[text(digest.as_str())])?
1264 .first(None)
1265 .await
1266 }
1267
1268 /// Every version of a workspace's packages of an ecosystem, oldest
1269 /// first: what an index of the whole registry is made from.
1270 pub async fn ecosystem_versions(&self, workspace: &str, ecosystem: &str, limit: u32) -> Result<Vec<VersionRow>> {
1271 self.prepare(
1272 &format!(
1273 "SELECT {} FROM versions v JOIN packages p ON p.id = v.package_id
1274 WHERE p.workspace = ? AND p.ecosystem = ? AND p.workspace_deleted_at IS NULL AND p.deleted_at IS NULL AND v.deleted_at IS NULL
1275 ORDER BY v.published_at, v.id LIMIT {limit}",
1276 VERSION_COLUMNS.split(", ").map(|c| format!("v.{c}")).collect::<Vec<_>>().join(", ")
1277 ),
1278 &[text(workspace), text(ecosystem)],
1279 )?
1280 .all()
1281 .await?
1282 .results()
1283 }
1284
1285 /// A workspace's packages of an ecosystem, by name.
1286 pub async fn packages_of(&self, workspace: &str, ecosystem: &str, limit: u32) -> Result<Vec<PackageRow>> {
1287 self.prepare(
1288 &format!(
1289 "SELECT {PACKAGE_COLUMNS} FROM packages WHERE workspace = ? AND ecosystem = ? AND workspace_deleted_at IS NULL AND deleted_at IS NULL ORDER BY name LIMIT {limit}"
1290 ),
1291 &[text(workspace), text(ecosystem)],
1292 )?
1293 .all()
1294 .await?
1295 .results()
1296 }
1297
1298 pub async fn package_by_id(&self, package_id: &str) -> Result<Option<PackageRow>> {
1299 self.prepare(&format!("SELECT {PACKAGE_COLUMNS} FROM packages WHERE id = ?"), &[text(package_id)])?
1300 .first(None)
1301 .await
1302 }
1303
1304 /// The files a workspace's packages of an ecosystem keep under `name`,
1305 /// newest first: how the NuGet symbol server finds a PDB.
1306 pub async fn files_named(&self, workspace: &str, ecosystem: &str, name: &str, limit: u32) -> Result<Vec<NamedFile>> {
1307 self.prepare(
1308 &format!(
1309 "SELECT v.package_id, f.digest, f.size FROM version_files f
1310 JOIN versions v ON v.id = f.version_id JOIN packages p ON p.id = v.package_id
1311 WHERE f.name = ? AND p.workspace = ? AND p.ecosystem = ? AND p.workspace_deleted_at IS NULL AND p.deleted_at IS NULL AND v.deleted_at IS NULL
1312 ORDER BY v.published_at DESC LIMIT {limit}"
1313 ),
1314 &[text(name), text(workspace), text(ecosystem)],
1315 )?
1316 .all()
1317 .await?
1318 .results()
1319 }
1320
1321 /// A workspace's Maven artifacts of one groupId (`com.acme:*`), by
1322 /// name: those named from `com.acme:` up to `com.acme;`, the
1323 /// character after `:`.
1324 pub async fn maven_group(&self, workspace: &str, group: &str, limit: u32) -> Result<Vec<PackageRow>> {
1325 self.prepare(
1326 &format!(
1327 "SELECT {PACKAGE_COLUMNS} FROM packages WHERE workspace = ? AND ecosystem = 'maven' AND name >= ? AND name < ?
1328 AND workspace_deleted_at IS NULL AND deleted_at IS NULL ORDER BY name LIMIT {limit}"
1329 ),
1330 &[text(workspace), text(&format!("{group}:")), text(&format!("{group};"))],
1331 )?
1332 .all()
1333 .await?
1334 .results()
1335 }
1336
1337 pub async fn forget_blob(&self, digest: &str) -> Result<()> {
1338 let d = [text(digest)];
1339 self.db
1340 .batch(vec![
1341 self.prepare("DELETE FROM checksums WHERE digest = ?", &d)?,
1342 self.prepare("DELETE FROM package_blobs WHERE digest = ?", &d)?,
1343 self.prepare("DELETE FROM workspace_blobs WHERE digest = ?", &d)?,
1344 self.prepare("DELETE FROM blobs WHERE digest = ?", &d)?,
1345 ])
1346 .await?;
1347 Ok(())
1348 }
1349
1350 /// Workspaces with packages linked to the repository.
1351 pub async fn workspaces_linked_to(&self, repo_id: &str) -> Result<Vec<String>> {
1352 #[derive(Deserialize)]
1353 struct Ws {
1354 workspace: String,
1355 }
1356 let rows: Vec<Ws> = self
1357 .prepare("SELECT DISTINCT workspace FROM packages WHERE repo_id = ?", &[text(repo_id)])?
1358 .all()
1359 .await?
1360 .results()?;
1361 Ok(rows.into_iter().map(|row| row.workspace).collect())
1362 }
1363
1364 pub async fn follow_visibility(&self, repo_id: &str, private: bool) -> Result<()> {
1365 self.prepare(
1366 "UPDATE packages SET visibility = ? WHERE repo_id = ?",
1367 &[text(if private { "private" } else { "public" }), text(repo_id)],
1368 )?
1369 .run()
1370 .await?;
1371 Ok(())
1372 }
1373
1374 /// Whether the package has ever had a version that is still kept,
1375 /// deleted ones included.
1376 pub async fn has_any_version(&self, package_id: &str) -> Result<bool> {
1377 let row: Option<serde_json::Value> = self
1378 .prepare("SELECT 1 AS found FROM versions WHERE package_id = ? LIMIT 1", &[text(package_id)])?
1379 .first(None)
1380 .await?;
1381 Ok(row.is_some())
1382 }
1383
1384 /// A deleted version by its id, version or digest: what a restore names,
1385 /// and what keeps a version string from being published again.
1386 pub async fn deleted_version(&self, package_id: &str, reference: &str) -> Result<Option<VersionRow>> {
1387 self.prepare(
1388 &format!(
1389 "SELECT {VERSION_COLUMNS} FROM versions WHERE package_id = ?1 AND deleted_at IS NOT NULL
1390 AND (id = ?2 OR version = ?2 OR digest = ?2) ORDER BY deleted_at DESC LIMIT 1"
1391 ),
1392 &[text(package_id), text(reference)],
1393 )?
1394 .first(None)
1395 .await
1396 }
1397
1398 /// A package's deleted versions, newest deletion first.
1399 pub async fn deleted_versions(&self, package_id: &str, limit: u32) -> Result<Vec<VersionRow>> {
1400 self.prepare(
1401 &format!(
1402 "SELECT {VERSION_COLUMNS} FROM versions WHERE package_id = ?1 AND deleted_at IS NOT NULL
1403 ORDER BY deleted_at DESC, id DESC LIMIT {limit}"
1404 ),
1405 &[text(package_id)],
1406 )?
1407 .all()
1408 .await?
1409 .results()
1410 }
1411
1412 /// Hides a version until it is restored or purged. Its files and tags
1413 /// stay, so a restore brings it back whole.
1414 pub async fn soft_delete_version(&self, version_id: &str, by: &str, now_ms: u64) -> Result<()> {
1415 self.prepare(
1416 "UPDATE versions SET deleted_at = ?, deleted_by = ? WHERE id = ? AND deleted_at IS NULL",
1417 &[text(&rfc3339(now_ms)), text(by), text(version_id)],
1418 )?
1419 .run()
1420 .await?;
1421 Ok(())
1422 }
1423
1424 pub async fn restore_version(&self, package_id: &str, version_id: &str, now_ms: u64) -> Result<()> {
1425 let now = rfc3339(now_ms);
1426 self.db
1427 .batch(vec![
1428 self.prepare("UPDATE versions SET deleted_at = NULL, deleted_by = NULL WHERE id = ?", &[text(version_id)])?,
1429 self.prepare("UPDATE packages SET updated_at = ? WHERE id = ?", &[text(&now), text(package_id)])?,
1430 ])
1431 .await?;
1432 Ok(())
1433 }
1434
1435 /// Hides a package and everything in it until it is restored or purged;
1436 /// its name stays taken meanwhile.
1437 pub async fn soft_delete_package(&self, package_id: &str, by: &str, now_ms: u64) -> Result<()> {
1438 self.db
1439 .batch(vec![
1440 self.prepare(
1441 "UPDATE packages SET deleted_at = ?, deleted_by = ? WHERE id = ? AND deleted_at IS NULL",
1442 &[text(&rfc3339(now_ms)), text(by), text(package_id)],
1443 )?,
1444 // Uploads in progress are not finished into a deleted package.
1445 self.prepare("DELETE FROM uploads WHERE package_id = ?", &[text(package_id)])?,
1446 ])
1447 .await?;
1448 Ok(())
1449 }
1450
1451 pub async fn restore_package(&self, package_id: &str, now_ms: u64) -> Result<()> {
1452 self.prepare(
1453 "UPDATE packages SET deleted_at = NULL, deleted_by = NULL, updated_at = ? WHERE id = ?",
1454 &[text(&rfc3339(now_ms)), text(package_id)],
1455 )?
1456 .run()
1457 .await?;
1458 Ok(())
1459 }
1460
1461 /// Packages deleted before `before`, for the purge.
1462 pub async fn expired_packages(&self, before: &str, limit: u32) -> Result<Vec<PackageRow>> {
1463 self.prepare(
1464 &format!("SELECT {PACKAGE_COLUMNS} FROM packages WHERE deleted_at IS NOT NULL AND deleted_at < ? ORDER BY deleted_at LIMIT {limit}"),
1465 &[text(before)],
1466 )?
1467 .all()
1468 .await?
1469 .results()
1470 }
1471
1472 /// Versions deleted before `before`, for the purge.
1473 pub async fn expired_versions(&self, before: &str, limit: u32) -> Result<Vec<VersionRow>> {
1474 self.prepare(
1475 &format!("SELECT {VERSION_COLUMNS} FROM versions WHERE deleted_at IS NOT NULL AND deleted_at < ? ORDER BY deleted_at LIMIT {limit}"),
1476 &[text(before)],
1477 )?
1478 .all()
1479 .await?
1480 .results()
1481 }
1482
1483 pub async fn set_inherit(&self, package_id: &str, inherit: bool, now_ms: u64) -> Result<()> {
1484 self.prepare(
1485 "UPDATE packages SET inherit_access = ?, updated_at = ? WHERE id = ?",
1486 &[num(u64::from(inherit)), text(&rfc3339(now_ms)), text(package_id)],
1487 )?
1488 .run()
1489 .await?;
1490 Ok(())
1491 }
1492
1493 /// The roles given on these packages, people first, then teams, by name.
1494 pub async fn access(&self, package_ids: &[String]) -> Result<Vec<AccessRow>> {
1495 let mut rows = Vec::new();
1496 for chunk in package_ids.chunks(50) {
1497 let marks = vec!["?"; chunk.len()].join(", ");
1498 let values: Vec<JsValue> = chunk.iter().map(|id| text(id)).collect();
1499 let found: Vec<AccessRow> = self
1500 .prepare(
1501 &format!(
1502 "SELECT package_id, grantee_kind, grantee_id, grantee_name, role, created_at FROM package_access
1503 WHERE package_id IN ({marks}) ORDER BY grantee_kind DESC, grantee_name"
1504 ),
1505 &values,
1506 )?
1507 .all()
1508 .await?
1509 .results()?;
1510 rows.extend(found);
1511 }
1512 Ok(rows)
1513 }
1514
1515 /// The repositories given access to these packages, by name.
1516 pub async fn actions_access(&self, package_ids: &[String]) -> Result<Vec<ActionsRow>> {
1517 let mut rows = Vec::new();
1518 for chunk in package_ids.chunks(50) {
1519 let marks = vec!["?"; chunk.len()].join(", ");
1520 let values: Vec<JsValue> = chunk.iter().map(|id| text(id)).collect();
1521 let found: Vec<ActionsRow> = self
1522 .prepare(
1523 &format!(
1524 "SELECT package_id, repo_id, repo_name, role, created_at FROM package_actions_access
1525 WHERE package_id IN ({marks}) ORDER BY repo_name"
1526 ),
1527 &values,
1528 )?
1529 .all()
1530 .await?
1531 .results()?;
1532 rows.extend(found);
1533 }
1534 Ok(rows)
1535 }
1536
1537 #[allow(clippy::too_many_arguments)]
1538 pub async fn set_access(&self, package_id: &str, kind: &str, id: &str, name: &str, role: &str, by: &str, now_ms: u64) -> Result<()> {
1539 self.prepare(
1540 "INSERT INTO package_access (package_id, grantee_kind, grantee_id, grantee_name, role, created_by, created_at)
1541 VALUES (?, ?, ?, ?, ?, ?, ?)
1542 ON CONFLICT (package_id, grantee_kind, grantee_id) DO UPDATE SET role = excluded.role, grantee_name = excluded.grantee_name",
1543 &[text(package_id), text(kind), text(id), text(name), text(role), text(by), text(&rfc3339(now_ms))],
1544 )?
1545 .run()
1546 .await?;
1547 Ok(())
1548 }
1549
1550 pub async fn remove_access(&self, package_id: &str, kind: &str, id: &str) -> Result<()> {
1551 self.prepare(
1552 "DELETE FROM package_access WHERE package_id = ? AND grantee_kind = ? AND grantee_id = ?",
1553 &[text(package_id), text(kind), text(id)],
1554 )?
1555 .run()
1556 .await?;
1557 Ok(())
1558 }
1559
1560 pub async fn set_actions(&self, package_id: &str, repo_id: &str, repo_name: &str, role: &str, by: &str, now_ms: u64) -> Result<()> {
1561 self.prepare(
1562 "INSERT INTO package_actions_access (package_id, repo_id, repo_name, role, created_by, created_at) VALUES (?, ?, ?, ?, ?, ?)
1563 ON CONFLICT (package_id, repo_id) DO UPDATE SET role = excluded.role, repo_name = excluded.repo_name",
1564 &[text(package_id), text(repo_id), text(repo_name), text(role), text(by), text(&rfc3339(now_ms))],
1565 )?
1566 .run()
1567 .await?;
1568 Ok(())
1569 }
1570
1571 pub async fn remove_actions(&self, package_id: &str, repo_id: &str) -> Result<()> {
1572 self.prepare(
1573 "DELETE FROM package_actions_access WHERE package_id = ? AND repo_id = ?",
1574 &[text(package_id), text(repo_id)],
1575 )?
1576 .run()
1577 .await?;
1578 Ok(())
1579 }
1580
1581 /// A team renamed (`slug` given) or deleted (`None`): its grants follow.
1582 pub async fn follow_team(&self, team_id: &str, slug: Option<&str>) -> Result<()> {
1583 let statement = match slug {
1584 Some(slug) => self.prepare(
1585 "UPDATE package_access SET grantee_name = ? WHERE grantee_kind = 'team' AND grantee_id = ?",
1586 &[text(slug), text(team_id)],
1587 )?,
1588 None => self.prepare("DELETE FROM package_access WHERE grantee_kind = 'team' AND grantee_id = ?", &[text(team_id)])?,
1589 };
1590 statement.run().await?;
1591 Ok(())
1592 }
1593}