g1t/services/packages/src/db.rs

978 lines38,658 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 /// Every version, newest published first, one per line: the summary
66 /// picks the highest of them (see `newest_version`).
67 pub latest_version: Option<String>,
68}
69
70/// The version a listing calls latest: the highest stable one by number
71/// (`v3.0.2` over `1.0.0`, whatever order they were published in, as an
72/// import publishes every tag at once), else the highest pre-release, else
73/// the newest published when none reads as a number.
74pub fn newest_version(versions: &str) -> Option<String> {
75 // Stable over pre-release, then by number, then pre-releases by label.
76 let parse = |version: &str| -> Option<(bool, Vec<u64>, String)> {
77 let bare = version.strip_prefix('v').unwrap_or(version);
78 let (core, pre) = match bare.split_once(['-', '+']) {
79 Some((core, rest)) if bare.as_bytes()[core.len()] == b'-' => (core, rest.to_owned()),
80 Some((core, _)) => (core, String::new()),
81 None => (bare, String::new()),
82 };
83 let parts = core.split('.').map(|part| part.parse::<u64>().ok()).collect::<Option<Vec<_>>>()?;
84 Some((pre.is_empty(), parts, pre))
85 };
86 let list: Vec<&str> = versions.lines().map(str::trim).filter(|v| !v.is_empty()).collect();
87 list.iter()
88 .filter_map(|v| parse(v).map(|key| (key, *v)))
89 .max_by(|a, b| a.0.cmp(&b.0))
90 .map(|(_, v)| v.to_owned())
91 .or_else(|| list.first().map(|v| (*v).to_owned()))
92}
93
94#[cfg(test)]
95mod newest_tests {
96 use super::newest_version;
97
98 #[test]
99 fn the_latest_is_the_highest_stable_version_not_the_last_published() {
100 assert_eq!(newest_version("1.0.0
1013.0.2
1022.0.0
1033.0.0").as_deref(), Some("3.0.2"));
104 assert_eq!(newest_version("v1.10.0
105v1.9.3").as_deref(), Some("v1.10.0"));
106 assert_eq!(newest_version("4.0.0-beta.1
1073.0.2").as_deref(), Some("3.0.2"));
108 assert_eq!(newest_version("4.0.0-beta.1
1094.0.0-alpha").as_deref(), Some("4.0.0-beta.1"));
110 assert_eq!(newest_version("dev-main
111nightly").as_deref(), Some("dev-main"));
112 assert_eq!(newest_version(""), None);
113 }
114}
115
116#[derive(Clone, Debug, Deserialize)]
117pub struct BlobRow {
118 pub digest: String,
119 pub size: u64,
120 pub media_type: Option<String>,
121 pub object_key: String,
122}
123
124#[derive(Clone, Debug, Deserialize)]
125pub struct VersionRow {
126 pub id: String,
127 pub package_id: String,
128 pub version: String,
129 pub digest: String,
130 pub size: u64,
131 pub metadata: String,
132 pub subject: Option<String>,
133 pub published_by: Option<String>,
134 pub published_at: String,
135 /// npm's deprecation message, when the version is deprecated.
136 #[serde(default)]
137 pub deprecated: Option<String>,
138}
139
140impl VersionRow {
141 pub fn meta(&self) -> serde_json::Value {
142 serde_json::from_str(&self.metadata).unwrap_or_default()
143 }
144
145 pub fn media_type(&self) -> Option<String> {
146 self.meta()["media_type"].as_str().map(str::to_owned)
147 }
148}
149
150#[derive(Clone, Debug, Deserialize)]
151pub struct TagRow {
152 pub tag: String,
153 pub version_id: String,
154 pub digest: String,
155 pub updated_at: String,
156}
157
158#[derive(Clone, Debug, Deserialize)]
159pub struct UploadRow {
160 pub id: String,
161 pub workspace: String,
162 pub package_id: String,
163 pub package: String,
164 pub multipart_id: Option<String>,
165 pub parts: String,
166 pub offset: u64,
167 pub tail: u64,
168 pub hash_state: String,
169}
170
171impl UploadRow {
172 pub fn progress(&self) -> Option<Progress> {
173 Some(Progress {
174 id: self.id.clone(),
175 multipart_id: self.multipart_id.clone(),
176 parts: serde_json::from_str(&self.parts).ok()?,
177 offset: self.offset,
178 tail: self.tail,
179 hasher: crate::digest::Sha256::restore(&self.hash_state)?,
180 })
181 }
182}
183
184/// One file of a version, as it is recorded.
185pub struct NewFile {
186 pub name: String,
187 pub digest: String,
188 pub size: u64,
189 pub media_type: Option<String>,
190}
191
192/// A version to record.
193pub struct NewVersion {
194 pub id: String,
195 pub package_id: String,
196 pub version: String,
197 pub digest: String,
198 pub size: u64,
199 pub metadata: String,
200 pub subject: Option<String>,
201 pub published_by: Option<String>,
202 pub files: Vec<NewFile>,
203}
204
205const PACKAGE_COLUMNS: &str =
206 "id, workspace, ecosystem, name, repo_id, repo_name, visibility, description, created_by, created_at, updated_at, downloads, workspace_deleted_at";
207/// Workspaces that are deleted, waiting to be purged or restored.
208const DELETED_WORKSPACES: &str = "SELECT workspace FROM packages WHERE workspace_deleted_at IS NOT NULL";
209const VERSION_COLUMNS: &str = "id, package_id, version, digest, size, metadata, subject, published_by, published_at, deprecated";
210
211pub struct Db {
212 pub db: D1Database,
213}
214
215impl Db {
216 fn prepare(&self, sql: &str, values: &[JsValue]) -> Result<D1PreparedStatement> {
217 self.db.prepare(sql).bind(values)
218 }
219
220 pub async fn package(&self, workspace: &str, ecosystem: &str, name: &str) -> Result<Option<PackageRow>> {
221 self.prepare(
222 &format!("SELECT {PACKAGE_COLUMNS} FROM packages WHERE workspace = ? AND ecosystem = ? AND name = ?"),
223 &[text(workspace), text(ecosystem), text(name)],
224 )?
225 .first(None)
226 .await
227 }
228
229 /// Makes a package unless one of the name is already there (a push
230 /// beside this one may have made it), and answers with the one kept.
231 #[allow(clippy::too_many_arguments)]
232 pub async fn create_package(
233 &self,
234 id: &str,
235 workspace: &str,
236 ecosystem: &str,
237 name: &str,
238 repo: Option<(&str, &str, bool)>,
239 created_by: &str,
240 now_ms: u64,
241 ) -> Result<PackageRow> {
242 let now = rfc3339(now_ms);
243 let visibility = match repo {
244 Some((_, _, private)) if !private => "public",
245 _ => "private",
246 };
247 self.prepare(
248 "INSERT OR IGNORE INTO packages (id, workspace, ecosystem, name, repo_id, repo_name, visibility, created_by, created_at, updated_at)
249 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
250 &[
251 text(id),
252 text(workspace),
253 text(ecosystem),
254 text(name),
255 opt(repo.map(|r| r.0)),
256 opt(repo.map(|r| r.1)),
257 text(visibility),
258 text(created_by),
259 text(&now),
260 text(&now),
261 ],
262 )?
263 .run()
264 .await?;
265 self.package(workspace, ecosystem, name)
266 .await?
267 .ok_or_else(|| worker::Error::RustError(format!("package {workspace}/{name} was not made")))
268 }
269
270 /// A workspace's packages with their counts, newest first.
271 pub async fn list(
272 &self,
273 workspace: &str,
274 ecosystem: Option<&str>,
275 repo_id: Option<&str>,
276 query: Option<&str>,
277 limit: u32,
278 ) -> Result<Vec<ListedRow>> {
279 let mut sql = format!(
280 "SELECT {}, \
281 (SELECT COUNT(*) FROM versions v WHERE v.package_id = p.id) AS version_count, \
282 (SELECT COALESCE(SUM(b.size), 0) FROM blobs b WHERE b.digest IN \
283 (SELECT vf.digest FROM version_files vf JOIN versions v ON v.id = vf.version_id WHERE v.package_id = p.id)) AS bytes, \
284 (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, \
285 (SELECT GROUP_CONCAT(version, char(10)) FROM (SELECT v.version FROM versions v WHERE v.package_id = p.id ORDER BY v.published_at DESC)) AS latest_version \
286 FROM packages p WHERE p.workspace = ? AND p.workspace_deleted_at IS NULL",
287 PACKAGE_COLUMNS.split(", ").map(|c| format!("p.{c}")).collect::<Vec<_>>().join(", ")
288 );
289 let mut values = vec![text(workspace)];
290 if let Some(ecosystem) = ecosystem {
291 sql.push_str(" AND p.ecosystem = ?");
292 values.push(text(ecosystem));
293 }
294 if let Some(repo_id) = repo_id {
295 sql.push_str(" AND p.repo_id = ?");
296 values.push(text(repo_id));
297 }
298 if let Some(query) = query.map(str::trim).filter(|q| !q.is_empty()) {
299 sql.push_str(" AND p.name LIKE ? ESCAPE '\\'");
300 let escaped = query.replace('\\', "\\\\").replace('%', "\\%").replace('_', "\\_");
301 values.push(text(&format!("%{}%", escaped.to_lowercase())));
302 }
303 sql.push_str(&format!(" ORDER BY p.updated_at DESC LIMIT {limit}"));
304 self.prepare(&sql, &values)?.all().await?.results()
305 }
306
307 pub async fn set_visibility(&self, package_id: &str, visibility: &str, now_ms: u64) -> Result<()> {
308 self.prepare(
309 "UPDATE packages SET visibility = ?, updated_at = ? WHERE id = ?",
310 &[text(visibility), text(&rfc3339(now_ms)), text(package_id)],
311 )?
312 .run()
313 .await?;
314 Ok(())
315 }
316
317 pub async fn set_link(&self, package_id: &str, repo: Option<(&str, &str)>, visibility: &str, now_ms: u64) -> Result<()> {
318 self.prepare(
319 "UPDATE packages SET repo_id = ?, repo_name = ?, visibility = ?, updated_at = ? WHERE id = ?",
320 &[opt(repo.map(|r| r.0)), opt(repo.map(|r| r.1)), text(visibility), text(&rfc3339(now_ms)), text(package_id)],
321 )?
322 .run()
323 .await?;
324 Ok(())
325 }
326
327 pub async fn add_downloads(&self, counts: &[(String, u64)]) -> Result<()> {
328 if counts.is_empty() {
329 return Ok(());
330 }
331 let mut batch = Vec::with_capacity(counts.len());
332 for (id, count) in counts {
333 batch.push(self.prepare("UPDATE packages SET downloads = downloads + ? WHERE id = ?", &[num(*count), text(id)])?);
334 }
335 self.db.batch(batch).await?;
336 Ok(())
337 }
338
339 /// A package and everything it holds. Its blobs are left to the sweep.
340 pub async fn delete_package(&self, package_id: &str) -> Result<()> {
341 let id = [text(package_id)];
342 self.db
343 .batch(vec![
344 self.prepare("DELETE FROM tags WHERE package_id = ?", &id)?,
345 self.prepare("DELETE FROM version_files WHERE version_id IN (SELECT id FROM versions WHERE package_id = ?)", &id)?,
346 self.prepare("DELETE FROM versions WHERE package_id = ?", &id)?,
347 self.prepare("DELETE FROM package_blobs WHERE package_id = ?", &id)?,
348 self.prepare("DELETE FROM uploads WHERE package_id = ?", &id)?,
349 self.prepare("DELETE FROM packages WHERE id = ?", &id)?,
350 ])
351 .await?;
352 Ok(())
353 }
354
355 pub async fn blob(&self, digest: &Digest) -> Result<Option<BlobRow>> {
356 self.prepare("SELECT digest, size, media_type, object_key FROM blobs WHERE digest = ?", &[text(digest.as_str())])?
357 .first(None)
358 .await
359 }
360
361 /// The blob, if this package may serve it.
362 pub async fn package_blob(&self, package_id: &str, digest: &Digest) -> Result<Option<BlobRow>> {
363 self.prepare(
364 "SELECT b.digest, b.size, b.media_type, b.object_key FROM package_blobs pb JOIN blobs b ON b.digest = pb.digest
365 WHERE pb.package_id = ? AND pb.digest = ?",
366 &[text(package_id), text(digest.as_str())],
367 )?
368 .first(None)
369 .await
370 }
371
372 /// Records a stored blob, or where it is now, and lets the package
373 /// serve it. Either way the blob is marked as just used.
374 pub async fn keep_blob(&self, package_id: &str, digest: &Digest, size: u64, media_type: Option<&str>, key: &str, now_ms: u64) -> Result<()> {
375 let now = rfc3339(now_ms);
376 self.db
377 .batch(vec![
378 self.prepare(
379 "INSERT INTO blobs (digest, size, media_type, object_key, created_at, touched_at) VALUES (?, ?, ?, ?, ?, ?)
380 ON CONFLICT (digest) DO UPDATE SET object_key = excluded.object_key, size = excluded.size",
381 &[text(digest.as_str()), num(size), opt(media_type), text(key), text(&now), text(&now)],
382 )?,
383 self.prepare("UPDATE blobs SET touched_at = ? WHERE digest = ?", &[text(&now), text(digest.as_str())])?,
384 self.prepare(
385 "INSERT OR IGNORE INTO package_blobs (package_id, digest, created_at) VALUES (?, ?, ?)",
386 &[text(package_id), text(digest.as_str()), text(&now)],
387 )?,
388 ])
389 .await?;
390 Ok(())
391 }
392
393 /// Lets the package serve a blob that is already stored.
394 pub async fn link_blob(&self, package_id: &str, digest: &Digest, now_ms: u64) -> Result<()> {
395 let now = rfc3339(now_ms);
396 self.db
397 .batch(vec![
398 self.prepare("UPDATE blobs SET touched_at = ? WHERE digest = ?", &[text(&now), text(digest.as_str())])?,
399 self.prepare(
400 "INSERT OR IGNORE INTO package_blobs (package_id, digest, created_at) VALUES (?, ?, ?)",
401 &[text(package_id), text(digest.as_str()), text(&now)],
402 )?,
403 ])
404 .await?;
405 Ok(())
406 }
407
408 /// Whether any version of the package names the blob.
409 pub async fn blob_in_use(&self, package_id: &str, digest: &Digest) -> Result<bool> {
410 let row: Option<serde_json::Value> = self
411 .prepare(
412 "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",
413 &[text(package_id), text(digest.as_str())],
414 )?
415 .first(None)
416 .await?;
417 Ok(row.is_some())
418 }
419
420 pub async fn unlink_blob(&self, package_id: &str, digest: &Digest) -> Result<()> {
421 self.prepare("DELETE FROM package_blobs WHERE package_id = ? AND digest = ?", &[text(package_id), text(digest.as_str())])?
422 .run()
423 .await?;
424 Ok(())
425 }
426
427 pub async fn create_upload(&self, row: &UploadRow, now_ms: u64) -> Result<()> {
428 self.prepare(
429 "INSERT INTO uploads (id, workspace, package_id, package, parts, \"offset\", tail, hash_state, created_at, expires_at)
430 VALUES (?, ?, ?, ?, '[]', 0, 0, ?, ?, ?)",
431 &[
432 text(&row.id),
433 text(&row.workspace),
434 text(&row.package_id),
435 text(&row.package),
436 text(&row.hash_state),
437 text(&rfc3339(now_ms)),
438 text(&rfc3339(now_ms + DAY_MS)),
439 ],
440 )?
441 .run()
442 .await?;
443 Ok(())
444 }
445
446 pub async fn upload(&self, id: &str) -> Result<Option<UploadRow>> {
447 self.prepare(
448 "SELECT id, workspace, package_id, package, multipart_id, parts, \"offset\" AS offset, tail, hash_state FROM uploads WHERE id = ?",
449 &[text(id)],
450 )?
451 .first(None)
452 .await
453 }
454
455 pub async fn save_progress(&self, progress: &Progress, now_ms: u64) -> Result<()> {
456 self.prepare(
457 "UPDATE uploads SET multipart_id = ?, parts = ?, \"offset\" = ?, tail = ?, hash_state = ?, expires_at = ? WHERE id = ?",
458 &[
459 opt(progress.multipart_id.as_deref()),
460 text(&serde_json::to_string(&progress.parts)?),
461 num(progress.offset),
462 num(progress.tail),
463 text(&progress.hasher.save()),
464 text(&rfc3339(now_ms + DAY_MS)),
465 text(&progress.id),
466 ],
467 )?
468 .run()
469 .await?;
470 Ok(())
471 }
472
473 pub async fn delete_upload(&self, id: &str) -> Result<()> {
474 self.prepare("DELETE FROM uploads WHERE id = ?", &[text(id)])?.run().await?;
475 Ok(())
476 }
477
478 pub async fn expired_uploads(&self, now_ms: u64, limit: u32) -> Result<Vec<UploadRow>> {
479 self.prepare(
480 &format!(
481 "SELECT id, workspace, package_id, package, multipart_id, parts, \"offset\" AS offset, tail, hash_state
482 FROM uploads WHERE expires_at < ? ORDER BY expires_at LIMIT {limit}"
483 ),
484 &[text(&rfc3339(now_ms))],
485 )?
486 .all()
487 .await?
488 .results()
489 }
490
491 pub async fn version_by_digest(&self, package_id: &str, digest: &str) -> Result<Option<VersionRow>> {
492 self.prepare(
493 &format!("SELECT {VERSION_COLUMNS} FROM versions WHERE package_id = ? AND digest = ? LIMIT 1"),
494 &[text(package_id), text(digest)],
495 )?
496 .first(None)
497 .await
498 }
499
500 pub async fn version_by_tag(&self, package_id: &str, tag: &str) -> Result<Option<VersionRow>> {
501 self.prepare(
502 &format!(
503 "SELECT {} FROM tags t JOIN versions v ON v.id = t.version_id WHERE t.package_id = ? AND t.tag = ?",
504 VERSION_COLUMNS.split(", ").map(|c| format!("v.{c}")).collect::<Vec<_>>().join(", ")
505 ),
506 &[text(package_id), text(tag)],
507 )?
508 .first(None)
509 .await
510 }
511
512 /// A version by its version, digest, or a tag that points to it.
513 pub async fn find_version(&self, package_id: &str, reference: &str) -> Result<Option<VersionRow>> {
514 if let Some(found) = self.version_by_digest(package_id, reference).await? {
515 return Ok(Some(found));
516 }
517 let by_version: Option<VersionRow> = self
518 .prepare(
519 &format!("SELECT {VERSION_COLUMNS} FROM versions WHERE package_id = ? AND version = ?"),
520 &[text(package_id), text(reference)],
521 )?
522 .first(None)
523 .await?;
524 if by_version.is_some() {
525 return Ok(by_version);
526 }
527 self.version_by_tag(package_id, reference).await
528 }
529
530 pub async fn versions(&self, package_id: &str, limit: u32) -> Result<Vec<VersionRow>> {
531 self.prepare(
532 &format!("SELECT {VERSION_COLUMNS} FROM versions WHERE package_id = ? ORDER BY published_at DESC, id DESC LIMIT {limit}"),
533 &[text(package_id)],
534 )?
535 .all()
536 .await?
537 .results()
538 }
539
540 /// Records a version unless one of its digest is there, moves `tag` to
541 /// it, and says whether anything was published: a new version, or a
542 /// tag that now points somewhere else.
543 pub async fn publish(&self, version: NewVersion, tag: Option<&str>, now_ms: u64) -> Result<(VersionRow, bool)> {
544 let now = rfc3339(now_ms);
545 let existing = self.version_by_digest(&version.package_id, &version.digest).await?;
546 let mut changed = existing.is_none();
547 let mut batch = Vec::new();
548 let version_id = match &existing {
549 Some(row) => row.id.clone(),
550 None => {
551 batch.push(self.prepare(
552 "INSERT INTO versions (id, package_id, version, digest, size, metadata, subject, published_by, published_at)
553 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)",
554 &[
555 text(&version.id),
556 text(&version.package_id),
557 text(&version.version),
558 text(&version.digest),
559 num(version.size),
560 text(&version.metadata),
561 opt(version.subject.as_deref()),
562 opt(version.published_by.as_deref()),
563 text(&now),
564 ],
565 )?);
566 for file in &version.files {
567 batch.push(self.prepare(
568 "INSERT OR IGNORE INTO version_files (version_id, name, digest, size, media_type) VALUES (?, ?, ?, ?, ?)",
569 &[text(&version.id), text(&file.name), text(&file.digest), num(file.size), opt(file.media_type.as_deref())],
570 )?);
571 }
572 version.id.clone()
573 }
574 };
575 if let Some(tag) = tag {
576 let before = self.version_by_tag(&version.package_id, tag).await?;
577 changed |= before.map(|row| row.id) != Some(version_id.clone());
578 batch.push(self.prepare(
579 "INSERT INTO tags (package_id, tag, version_id, updated_at) VALUES (?, ?, ?, ?)
580 ON CONFLICT (package_id, tag) DO UPDATE SET version_id = excluded.version_id, updated_at = excluded.updated_at",
581 &[text(&version.package_id), text(tag), text(&version_id), text(&now)],
582 )?);
583 }
584 if changed {
585 batch.push(self.prepare("UPDATE packages SET updated_at = ? WHERE id = ?", &[text(&now), text(&version.package_id)])?);
586 }
587 if !batch.is_empty() {
588 self.db.batch(batch).await?;
589 }
590 let row = self
591 .version_by_digest(&version.package_id, &version.digest)
592 .await?
593 .ok_or_else(|| worker::Error::RustError("the version was not recorded".into()))?;
594 Ok((row, changed))
595 }
596
597 pub async fn tags(&self, package_id: &str) -> Result<Vec<TagRow>> {
598 self.prepare(
599 "SELECT t.tag, t.version_id, v.digest, t.updated_at FROM tags t JOIN versions v ON v.id = t.version_id
600 WHERE t.package_id = ? ORDER BY t.tag",
601 &[text(package_id)],
602 )?
603 .all()
604 .await?
605 .results()
606 }
607
608 /// Tag names in order, `n` of them after `last`.
609 pub async fn tag_names(&self, package_id: &str, last: Option<&str>, n: u32) -> Result<Vec<String>> {
610 #[derive(Deserialize)]
611 struct Name {
612 tag: String,
613 }
614 let rows: Vec<Name> = self
615 .prepare(
616 &format!("SELECT tag FROM tags WHERE package_id = ? AND tag > ? ORDER BY tag LIMIT {n}"),
617 &[text(package_id), text(last.unwrap_or(""))],
618 )?
619 .all()
620 .await?
621 .results()?;
622 Ok(rows.into_iter().map(|row| row.tag).collect())
623 }
624
625 pub async fn referrers(&self, package_id: &str, subject: &str) -> Result<Vec<VersionRow>> {
626 self.prepare(
627 &format!("SELECT {VERSION_COLUMNS} FROM versions WHERE package_id = ? AND subject = ? ORDER BY published_at"),
628 &[text(package_id), text(subject)],
629 )?
630 .all()
631 .await?
632 .results()
633 }
634
635 /// The package of an ecosystem built from a repository (Composer's).
636 pub async fn package_for_repo(&self, repo_id: &str, ecosystem: &str) -> Result<Option<PackageRow>> {
637 self.prepare(
638 &format!("SELECT {PACKAGE_COLUMNS} FROM packages WHERE repo_id = ? AND ecosystem = ? LIMIT 1"),
639 &[text(repo_id), text(ecosystem)],
640 )?
641 .first(None)
642 .await
643 }
644
645 pub async fn rename_package(&self, package_id: &str, name: &str, now_ms: u64) -> Result<()> {
646 self.prepare(
647 "UPDATE packages SET name = ?, updated_at = ? WHERE id = ?",
648 &[text(name), text(&rfc3339(now_ms)), text(package_id)],
649 )?
650 .run()
651 .await?;
652 Ok(())
653 }
654
655 /// Records a version, in place of one of the same version string
656 /// (a tag or branch that moved): its files and tags go with the old one.
657 pub async fn replace_version(&self, version: &NewVersion, now_ms: u64) -> Result<()> {
658 let now = rfc3339(now_ms);
659 let old = [text(&version.package_id), text(&version.version)];
660 let old_ids = "SELECT id FROM versions WHERE package_id = ? AND version = ?";
661 let mut batch = vec![
662 self.prepare(&format!("DELETE FROM tags WHERE version_id IN ({old_ids})"), &old)?,
663 self.prepare(&format!("DELETE FROM version_files WHERE version_id IN ({old_ids})"), &old)?,
664 self.prepare("DELETE FROM versions WHERE package_id = ? AND version = ?", &old)?,
665 self.prepare(
666 "INSERT INTO versions (id, package_id, version, digest, size, metadata, subject, published_by, published_at)
667 VALUES (?, ?, ?, ?, ?, ?, NULL, ?, ?)",
668 &[
669 text(&version.id),
670 text(&version.package_id),
671 text(&version.version),
672 text(&version.digest),
673 num(version.size),
674 text(&version.metadata),
675 opt(version.published_by.as_deref()),
676 text(&now),
677 ],
678 )?,
679 ];
680 for file in &version.files {
681 batch.push(self.prepare(
682 "INSERT OR IGNORE INTO version_files (version_id, name, digest, size, media_type) VALUES (?, ?, ?, ?, ?)",
683 &[text(&version.id), text(&file.name), text(&file.digest), num(file.size), opt(file.media_type.as_deref())],
684 )?);
685 }
686 batch.push(self.prepare("UPDATE packages SET updated_at = ? WHERE id = ?", &[text(&now), text(&version.package_id)])?);
687 self.db.batch(batch).await?;
688 Ok(())
689 }
690
691 /// The zip already made of a commit of the package, if one was.
692 pub async fn dist_for_commit(&self, package_id: &str, commit: &str) -> Result<Option<BlobRow>> {
693 self.prepare(
694 "SELECT b.digest, b.size, b.media_type, b.object_key FROM versions v
695 JOIN version_files vf ON vf.version_id = v.id AND vf.name = 'dist'
696 JOIN blobs b ON b.digest = vf.digest
697 WHERE v.package_id = ? AND v.digest = ? LIMIT 1",
698 &[text(package_id), text(commit)],
699 )?
700 .first(None)
701 .await
702 }
703
704 /// Records the zip of a commit as a file of every version at it.
705 pub async fn add_dist(&self, package_id: &str, commit: &str, digest: &Digest, size: u64) -> Result<()> {
706 let at = [text(package_id), text(commit)];
707 self.db
708 .batch(vec![
709 self.prepare(
710 "INSERT OR IGNORE INTO version_files (version_id, name, digest, size, media_type)
711 SELECT id, 'dist', ?, ?, 'application/zip' FROM versions WHERE package_id = ? AND digest = ?",
712 &[text(digest.as_str()), num(size), text(package_id), text(commit)],
713 )?,
714 self.prepare("UPDATE versions SET size = ? WHERE package_id = ? AND digest = ?", &[num(size), at[0].clone(), at[1].clone()])?,
715 ])
716 .await?;
717 Ok(())
718 }
719
720 /// Where the Composer backfill is: the last repository done, and
721 /// whether it went through them all.
722 pub async fn backfill(&self) -> Result<(Option<String>, bool)> {
723 #[derive(Deserialize)]
724 struct Row {
725 after: Option<String>,
726 finished_at: Option<String>,
727 }
728 let row: Option<Row> = self
729 .prepare("SELECT after, finished_at FROM composer_backfill WHERE key = 'repos'", &[])?
730 .first(None)
731 .await?;
732 Ok(row.map_or((None, false), |row| (row.after, row.finished_at.is_some())))
733 }
734
735 pub async fn set_backfill(&self, after: Option<&str>, finished: bool, now_ms: u64) -> Result<()> {
736 self.prepare(
737 "INSERT INTO composer_backfill (key, after, finished_at) VALUES ('repos', ?, ?)
738 ON CONFLICT (key) DO UPDATE SET after = excluded.after, finished_at = excluded.finished_at",
739 &[opt(after), if finished { text(&rfc3339(now_ms)) } else { JsValue::NULL }],
740 )?
741 .run()
742 .await?;
743 Ok(())
744 }
745
746 /// Points `tag` at a version, made or moved.
747 pub async fn set_tag(&self, package_id: &str, tag: &str, version_id: &str, now_ms: u64) -> Result<()> {
748 self.prepare(
749 "INSERT INTO tags (package_id, tag, version_id, updated_at) VALUES (?, ?, ?, ?)
750 ON CONFLICT (package_id, tag) DO UPDATE SET version_id = excluded.version_id, updated_at = excluded.updated_at",
751 &[text(package_id), text(tag), text(version_id), text(&rfc3339(now_ms))],
752 )?
753 .run()
754 .await?;
755 Ok(())
756 }
757
758 /// A version by its version string alone.
759 pub async fn version_named(&self, package_id: &str, version: &str) -> Result<Option<VersionRow>> {
760 self.prepare(
761 &format!("SELECT {VERSION_COLUMNS} FROM versions WHERE package_id = ? AND version = ?"),
762 &[text(package_id), text(version)],
763 )?
764 .first(None)
765 .await
766 }
767
768 pub async fn set_deprecated(&self, version_id: &str, message: Option<&str>) -> Result<()> {
769 self.prepare("UPDATE versions SET deprecated = ? WHERE id = ?", &[opt(message), text(version_id)])?
770 .run()
771 .await?;
772 Ok(())
773 }
774
775 /// The README shown for the package, kept as a blob, and its description.
776 pub async fn set_readme(&self, package_id: &str, digest: Option<&str>, description: Option<&str>, now_ms: u64) -> Result<()> {
777 self.prepare(
778 "UPDATE packages SET readme_digest = ?, description = ?, updated_at = ? WHERE id = ?",
779 &[opt(digest), opt(description), text(&rfc3339(now_ms)), text(package_id)],
780 )?
781 .run()
782 .await?;
783 Ok(())
784 }
785
786 pub async fn readme_digest(&self, package_id: &str) -> Result<Option<String>> {
787 #[derive(Deserialize)]
788 struct Readme {
789 readme_digest: Option<String>,
790 }
791 let row: Option<Readme> = self
792 .prepare("SELECT readme_digest FROM packages WHERE id = ?", &[text(package_id)])?
793 .first(None)
794 .await?;
795 Ok(row.and_then(|r| r.readme_digest))
796 }
797
798 pub async fn touch_package(&self, package_id: &str, now_ms: u64) -> Result<()> {
799 self.prepare("UPDATE packages SET updated_at = ? WHERE id = ?", &[text(&rfc3339(now_ms)), text(package_id)])?
800 .run()
801 .await?;
802 Ok(())
803 }
804
805 pub async fn delete_tag(&self, package_id: &str, tag: &str) -> Result<()> {
806 self.prepare("DELETE FROM tags WHERE package_id = ? AND tag = ?", &[text(package_id), text(tag)])?
807 .run()
808 .await?;
809 Ok(())
810 }
811
812 /// A version, its tags and its files. Blobs are left to the sweep.
813 pub async fn delete_version(&self, version_id: &str) -> Result<()> {
814 let id = [text(version_id)];
815 self.db
816 .batch(vec![
817 self.prepare("DELETE FROM tags WHERE version_id = ?", &id)?,
818 self.prepare("DELETE FROM version_files WHERE version_id = ?", &id)?,
819 self.prepare("DELETE FROM versions WHERE id = ?", &id)?,
820 ])
821 .await?;
822 Ok(())
823 }
824
825 /// Works out again what the workspace stores: each blob its versions
826 /// use, once, public when any public package uses it.
827 pub async fn measure(&self, workspace: &str) -> Result<()> {
828 let ws = [text(workspace)];
829 self.db
830 .batch(vec![
831 self.prepare("DELETE FROM workspace_blobs WHERE workspace = ?", &ws)?,
832 self.prepare(
833 "INSERT INTO workspace_blobs (workspace, digest, size, public)
834 SELECT p.workspace, vf.digest, MAX(vf.size), MAX(p.visibility = 'public')
835 FROM version_files vf JOIN versions v ON v.id = vf.version_id JOIN packages p ON p.id = v.package_id
836 WHERE p.workspace = ? GROUP BY vf.digest",
837 &ws,
838 )?,
839 ])
840 .await?;
841 Ok(())
842 }
843
844 pub async fn storage(&self, workspace: &str) -> Result<(u64, u64)> {
845 #[derive(Deserialize)]
846 struct Sums {
847 public_bytes: Option<u64>,
848 private_bytes: Option<u64>,
849 }
850 let sums: Option<Sums> = self
851 .prepare(
852 &format!(
853 "SELECT SUM(CASE WHEN public = 1 THEN size ELSE 0 END) AS public_bytes,
854 SUM(CASE WHEN public = 1 THEN 0 ELSE size END) AS private_bytes
855 FROM workspace_blobs WHERE workspace = ? AND workspace NOT IN ({DELETED_WORKSPACES})"
856 ),
857 &[text(workspace)],
858 )?
859 .first(None)
860 .await?;
861 Ok(sums.map_or((0, 0), |s| (s.public_bytes.unwrap_or(0), s.private_bytes.unwrap_or(0))))
862 }
863
864 /// Hides (`deleted_at` given) or shows again (`None`) a workspace's
865 /// packages. Hiding keeps a mark already set; showing clears only marks.
866 pub async fn mark_workspace(&self, workspace: &str, deleted_at: Option<&str>) -> Result<()> {
867 let statement = match deleted_at {
868 Some(at) => self.prepare(
869 "UPDATE packages SET workspace_deleted_at = ? WHERE workspace = ? AND workspace_deleted_at IS NULL",
870 &[text(at), text(workspace)],
871 )?,
872 None => self.prepare(
873 "UPDATE packages SET workspace_deleted_at = NULL WHERE workspace = ? AND workspace_deleted_at IS NOT NULL",
874 &[text(workspace)],
875 )?,
876 };
877 statement.run().await?;
878 Ok(())
879 }
880
881 /// Whether the workspace is deleted, as its packages say.
882 pub async fn workspace_hidden(&self, workspace: &str) -> Result<bool> {
883 let row: Option<serde_json::Value> = self
884 .prepare("SELECT 1 AS hidden FROM packages WHERE workspace = ? AND workspace_deleted_at IS NOT NULL LIMIT 1", &[text(workspace)])?
885 .first(None)
886 .await?;
887 Ok(row.is_some())
888 }
889
890 /// Every workspace's package storage, by workspace.
891 pub async fn storage_all(&self) -> Result<Vec<g1t_contracts::packages::WorkspacePackageStorage>> {
892 self.prepare(
893 &format!(
894 "SELECT workspace,
895 COALESCE(SUM(CASE WHEN public = 1 THEN size ELSE 0 END), 0) AS public_bytes,
896 COALESCE(SUM(CASE WHEN public = 1 THEN 0 ELSE size END), 0) AS private_bytes
897 FROM workspace_blobs WHERE workspace NOT IN ({DELETED_WORKSPACES}) GROUP BY workspace ORDER BY workspace"
898 ),
899 &[],
900 )?
901 .all()
902 .await?
903 .results()
904 }
905
906 /// Which of `digests` the workspace already holds.
907 pub async fn held(&self, workspace: &str, digests: &[String]) -> Result<std::collections::HashSet<String>> {
908 #[derive(Deserialize)]
909 struct Held {
910 digest: String,
911 }
912 let mut held = std::collections::HashSet::new();
913 for chunk in digests.chunks(50) {
914 let marks = vec!["?"; chunk.len()].join(", ");
915 let mut values = vec![text(workspace)];
916 values.extend(chunk.iter().map(|d| text(d)));
917 let rows: Vec<Held> = self
918 .prepare(&format!("SELECT digest FROM workspace_blobs WHERE workspace = ? AND digest IN ({marks})"), &values)?
919 .all()
920 .await?
921 .results()?;
922 held.extend(rows.into_iter().map(|row| row.digest));
923 }
924 Ok(held)
925 }
926
927 /// Blobs no version uses and nothing has touched for a day.
928 pub async fn unused_blobs(&self, now_ms: u64, limit: u32) -> Result<Vec<BlobRow>> {
929 self.prepare(
930 &format!(
931 "SELECT digest, size, media_type, object_key FROM blobs b
932 WHERE touched_at < ? AND NOT EXISTS (SELECT 1 FROM version_files vf WHERE vf.digest = b.digest)
933 AND NOT EXISTS (SELECT 1 FROM packages p WHERE p.readme_digest = b.digest)
934 ORDER BY touched_at LIMIT {limit}"
935 ),
936 &[text(&rfc3339(now_ms.saturating_sub(DAY_MS)))],
937 )?
938 .all()
939 .await?
940 .results()
941 }
942
943 pub async fn forget_blob(&self, digest: &str) -> Result<()> {
944 let d = [text(digest)];
945 self.db
946 .batch(vec![
947 self.prepare("DELETE FROM package_blobs WHERE digest = ?", &d)?,
948 self.prepare("DELETE FROM workspace_blobs WHERE digest = ?", &d)?,
949 self.prepare("DELETE FROM blobs WHERE digest = ?", &d)?,
950 ])
951 .await?;
952 Ok(())
953 }
954
955 /// Workspaces with packages linked to the repository.
956 pub async fn workspaces_linked_to(&self, repo_id: &str) -> Result<Vec<String>> {
957 #[derive(Deserialize)]
958 struct Ws {
959 workspace: String,
960 }
961 let rows: Vec<Ws> = self
962 .prepare("SELECT DISTINCT workspace FROM packages WHERE repo_id = ?", &[text(repo_id)])?
963 .all()
964 .await?
965 .results()?;
966 Ok(rows.into_iter().map(|row| row.workspace).collect())
967 }
968
969 pub async fn follow_visibility(&self, repo_id: &str, private: bool) -> Result<()> {
970 self.prepare(
971 "UPDATE packages SET visibility = ? WHERE repo_id = ?",
972 &[text(if private { "private" } else { "public" }), text(repo_id)],
973 )?
974 .run()
975 .await?;
976 Ok(())
977 }
978}