g1t/services/packages/src/db.rs

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