g1t/services/packages/src/db.rs

756 lines28,854 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 pub latest_version: Option<String>,
66}
67
68#[derive(Clone, Debug, Deserialize)]
69pub struct BlobRow {
70 pub digest: String,
71 pub size: u64,
72 pub media_type: Option<String>,
73 pub object_key: String,
74}
75
76#[derive(Clone, Debug, Deserialize)]
77pub struct VersionRow {
78 pub id: String,
79 pub package_id: String,
80 pub version: String,
81 pub digest: String,
82 pub size: u64,
83 pub metadata: String,
84 pub subject: Option<String>,
85 pub published_by: Option<String>,
86 pub published_at: String,
87}
88
89impl VersionRow {
90 pub fn meta(&self) -> serde_json::Value {
91 serde_json::from_str(&self.metadata).unwrap_or_default()
92 }
93
94 pub fn media_type(&self) -> Option<String> {
95 self.meta()["media_type"].as_str().map(str::to_owned)
96 }
97}
98
99#[derive(Clone, Debug, Deserialize)]
100pub struct TagRow {
101 pub tag: String,
102 pub version_id: String,
103 pub digest: String,
104 pub updated_at: String,
105}
106
107#[derive(Clone, Debug, Deserialize)]
108pub struct UploadRow {
109 pub id: String,
110 pub workspace: String,
111 pub package_id: String,
112 pub package: String,
113 pub multipart_id: Option<String>,
114 pub parts: String,
115 pub offset: u64,
116 pub tail: u64,
117 pub hash_state: String,
118}
119
120impl UploadRow {
121 pub fn progress(&self) -> Option<Progress> {
122 Some(Progress {
123 id: self.id.clone(),
124 multipart_id: self.multipart_id.clone(),
125 parts: serde_json::from_str(&self.parts).ok()?,
126 offset: self.offset,
127 tail: self.tail,
128 hasher: crate::digest::Sha256::restore(&self.hash_state)?,
129 })
130 }
131}
132
133/// One file of a version, as it is recorded.
134pub struct NewFile {
135 pub name: String,
136 pub digest: String,
137 pub size: u64,
138 pub media_type: Option<String>,
139}
140
141/// A version to record.
142pub struct NewVersion {
143 pub id: String,
144 pub package_id: String,
145 pub version: String,
146 pub digest: String,
147 pub size: u64,
148 pub metadata: String,
149 pub subject: Option<String>,
150 pub published_by: Option<String>,
151 pub files: Vec<NewFile>,
152}
153
154const PACKAGE_COLUMNS: &str =
155 "id, workspace, ecosystem, name, repo_id, repo_name, visibility, description, created_by, created_at, updated_at, downloads, workspace_deleted_at";
156/// Workspaces that are deleted, waiting to be purged or restored.
157const DELETED_WORKSPACES: &str = "SELECT workspace FROM packages WHERE workspace_deleted_at IS NOT NULL";
158const VERSION_COLUMNS: &str = "id, package_id, version, digest, size, metadata, subject, published_by, published_at";
159
160pub struct Db {
161 pub db: D1Database,
162}
163
164impl Db {
165 fn prepare(&self, sql: &str, values: &[JsValue]) -> Result<D1PreparedStatement> {
166 self.db.prepare(sql).bind(values)
167 }
168
169 pub async fn package(&self, workspace: &str, ecosystem: &str, name: &str) -> Result<Option<PackageRow>> {
170 self.prepare(
171 &format!("SELECT {PACKAGE_COLUMNS} FROM packages WHERE workspace = ? AND ecosystem = ? AND name = ?"),
172 &[text(workspace), text(ecosystem), text(name)],
173 )?
174 .first(None)
175 .await
176 }
177
178 /// Makes a package unless one of the name is already there (a push
179 /// beside this one may have made it), and answers with the one kept.
180 #[allow(clippy::too_many_arguments)]
181 pub async fn create_package(
182 &self,
183 id: &str,
184 workspace: &str,
185 ecosystem: &str,
186 name: &str,
187 repo: Option<(&str, &str, bool)>,
188 created_by: &str,
189 now_ms: u64,
190 ) -> Result<PackageRow> {
191 let now = rfc3339(now_ms);
192 let visibility = match repo {
193 Some((_, _, private)) if !private => "public",
194 _ => "private",
195 };
196 self.prepare(
197 "INSERT OR IGNORE INTO packages (id, workspace, ecosystem, name, repo_id, repo_name, visibility, created_by, created_at, updated_at)
198 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?, ?)",
199 &[
200 text(id),
201 text(workspace),
202 text(ecosystem),
203 text(name),
204 opt(repo.map(|r| r.0)),
205 opt(repo.map(|r| r.1)),
206 text(visibility),
207 text(created_by),
208 text(&now),
209 text(&now),
210 ],
211 )?
212 .run()
213 .await?;
214 self.package(workspace, ecosystem, name)
215 .await?
216 .ok_or_else(|| worker::Error::RustError(format!("package {workspace}/{name} was not made")))
217 }
218
219 /// A workspace's packages with their counts, newest first.
220 pub async fn list(
221 &self,
222 workspace: &str,
223 ecosystem: Option<&str>,
224 repo_id: Option<&str>,
225 query: Option<&str>,
226 limit: u32,
227 ) -> Result<Vec<ListedRow>> {
228 let mut sql = format!(
229 "SELECT {}, \
230 (SELECT COUNT(*) FROM versions v WHERE v.package_id = p.id) AS version_count, \
231 (SELECT COALESCE(SUM(b.size), 0) FROM blobs b WHERE b.digest IN \
232 (SELECT vf.digest FROM version_files vf JOIN versions v ON v.id = vf.version_id WHERE v.package_id = p.id)) AS bytes, \
233 (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, \
234 (SELECT v.version FROM versions v WHERE v.package_id = p.id ORDER BY v.published_at DESC LIMIT 1) AS latest_version \
235 FROM packages p WHERE p.workspace = ? AND p.workspace_deleted_at IS NULL",
236 PACKAGE_COLUMNS.split(", ").map(|c| format!("p.{c}")).collect::<Vec<_>>().join(", ")
237 );
238 let mut values = vec![text(workspace)];
239 if let Some(ecosystem) = ecosystem {
240 sql.push_str(" AND p.ecosystem = ?");
241 values.push(text(ecosystem));
242 }
243 if let Some(repo_id) = repo_id {
244 sql.push_str(" AND p.repo_id = ?");
245 values.push(text(repo_id));
246 }
247 if let Some(query) = query.map(str::trim).filter(|q| !q.is_empty()) {
248 sql.push_str(" AND p.name LIKE ? ESCAPE '\\'");
249 let escaped = query.replace('\\', "\\\\").replace('%', "\\%").replace('_', "\\_");
250 values.push(text(&format!("%{}%", escaped.to_lowercase())));
251 }
252 sql.push_str(&format!(" ORDER BY p.updated_at DESC LIMIT {limit}"));
253 self.prepare(&sql, &values)?.all().await?.results()
254 }
255
256 pub async fn set_visibility(&self, package_id: &str, visibility: &str, now_ms: u64) -> Result<()> {
257 self.prepare(
258 "UPDATE packages SET visibility = ?, updated_at = ? WHERE id = ?",
259 &[text(visibility), text(&rfc3339(now_ms)), text(package_id)],
260 )?
261 .run()
262 .await?;
263 Ok(())
264 }
265
266 pub async fn set_link(&self, package_id: &str, repo: Option<(&str, &str)>, visibility: &str, now_ms: u64) -> Result<()> {
267 self.prepare(
268 "UPDATE packages SET repo_id = ?, repo_name = ?, visibility = ?, updated_at = ? WHERE id = ?",
269 &[opt(repo.map(|r| r.0)), opt(repo.map(|r| r.1)), text(visibility), text(&rfc3339(now_ms)), text(package_id)],
270 )?
271 .run()
272 .await?;
273 Ok(())
274 }
275
276 pub async fn add_downloads(&self, counts: &[(String, u64)]) -> Result<()> {
277 if counts.is_empty() {
278 return Ok(());
279 }
280 let mut batch = Vec::with_capacity(counts.len());
281 for (id, count) in counts {
282 batch.push(self.prepare("UPDATE packages SET downloads = downloads + ? WHERE id = ?", &[num(*count), text(id)])?);
283 }
284 self.db.batch(batch).await?;
285 Ok(())
286 }
287
288 /// A package and everything it holds. Its blobs are left to the sweep.
289 pub async fn delete_package(&self, package_id: &str) -> Result<()> {
290 let id = [text(package_id)];
291 self.db
292 .batch(vec![
293 self.prepare("DELETE FROM tags WHERE package_id = ?", &id)?,
294 self.prepare("DELETE FROM version_files WHERE version_id IN (SELECT id FROM versions WHERE package_id = ?)", &id)?,
295 self.prepare("DELETE FROM versions WHERE package_id = ?", &id)?,
296 self.prepare("DELETE FROM package_blobs WHERE package_id = ?", &id)?,
297 self.prepare("DELETE FROM uploads WHERE package_id = ?", &id)?,
298 self.prepare("DELETE FROM packages WHERE id = ?", &id)?,
299 ])
300 .await?;
301 Ok(())
302 }
303
304 pub async fn blob(&self, digest: &Digest) -> Result<Option<BlobRow>> {
305 self.prepare("SELECT digest, size, media_type, object_key FROM blobs WHERE digest = ?", &[text(digest.as_str())])?
306 .first(None)
307 .await
308 }
309
310 /// The blob, if this package may serve it.
311 pub async fn package_blob(&self, package_id: &str, digest: &Digest) -> Result<Option<BlobRow>> {
312 self.prepare(
313 "SELECT b.digest, b.size, b.media_type, b.object_key FROM package_blobs pb JOIN blobs b ON b.digest = pb.digest
314 WHERE pb.package_id = ? AND pb.digest = ?",
315 &[text(package_id), text(digest.as_str())],
316 )?
317 .first(None)
318 .await
319 }
320
321 /// Records a stored blob, or where it is now, and lets the package
322 /// serve it. Either way the blob is marked as just used.
323 pub async fn keep_blob(&self, package_id: &str, digest: &Digest, size: u64, media_type: Option<&str>, key: &str, now_ms: u64) -> Result<()> {
324 let now = rfc3339(now_ms);
325 self.db
326 .batch(vec![
327 self.prepare(
328 "INSERT INTO blobs (digest, size, media_type, object_key, created_at, touched_at) VALUES (?, ?, ?, ?, ?, ?)
329 ON CONFLICT (digest) DO UPDATE SET object_key = excluded.object_key, size = excluded.size",
330 &[text(digest.as_str()), num(size), opt(media_type), text(key), text(&now), text(&now)],
331 )?,
332 self.prepare("UPDATE blobs SET touched_at = ? WHERE digest = ?", &[text(&now), text(digest.as_str())])?,
333 self.prepare(
334 "INSERT OR IGNORE INTO package_blobs (package_id, digest, created_at) VALUES (?, ?, ?)",
335 &[text(package_id), text(digest.as_str()), text(&now)],
336 )?,
337 ])
338 .await?;
339 Ok(())
340 }
341
342 /// Lets the package serve a blob that is already stored.
343 pub async fn link_blob(&self, package_id: &str, digest: &Digest, now_ms: u64) -> Result<()> {
344 let now = rfc3339(now_ms);
345 self.db
346 .batch(vec![
347 self.prepare("UPDATE blobs SET touched_at = ? WHERE digest = ?", &[text(&now), text(digest.as_str())])?,
348 self.prepare(
349 "INSERT OR IGNORE INTO package_blobs (package_id, digest, created_at) VALUES (?, ?, ?)",
350 &[text(package_id), text(digest.as_str()), text(&now)],
351 )?,
352 ])
353 .await?;
354 Ok(())
355 }
356
357 /// Whether any version of the package names the blob.
358 pub async fn blob_in_use(&self, package_id: &str, digest: &Digest) -> Result<bool> {
359 let row: Option<serde_json::Value> = self
360 .prepare(
361 "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",
362 &[text(package_id), text(digest.as_str())],
363 )?
364 .first(None)
365 .await?;
366 Ok(row.is_some())
367 }
368
369 pub async fn unlink_blob(&self, package_id: &str, digest: &Digest) -> Result<()> {
370 self.prepare("DELETE FROM package_blobs WHERE package_id = ? AND digest = ?", &[text(package_id), text(digest.as_str())])?
371 .run()
372 .await?;
373 Ok(())
374 }
375
376 pub async fn create_upload(&self, row: &UploadRow, now_ms: u64) -> Result<()> {
377 self.prepare(
378 "INSERT INTO uploads (id, workspace, package_id, package, parts, \"offset\", tail, hash_state, created_at, expires_at)
379 VALUES (?, ?, ?, ?, '[]', 0, 0, ?, ?, ?)",
380 &[
381 text(&row.id),
382 text(&row.workspace),
383 text(&row.package_id),
384 text(&row.package),
385 text(&row.hash_state),
386 text(&rfc3339(now_ms)),
387 text(&rfc3339(now_ms + DAY_MS)),
388 ],
389 )?
390 .run()
391 .await?;
392 Ok(())
393 }
394
395 pub async fn upload(&self, id: &str) -> Result<Option<UploadRow>> {
396 self.prepare(
397 "SELECT id, workspace, package_id, package, multipart_id, parts, \"offset\" AS offset, tail, hash_state FROM uploads WHERE id = ?",
398 &[text(id)],
399 )?
400 .first(None)
401 .await
402 }
403
404 pub async fn save_progress(&self, progress: &Progress, now_ms: u64) -> Result<()> {
405 self.prepare(
406 "UPDATE uploads SET multipart_id = ?, parts = ?, \"offset\" = ?, tail = ?, hash_state = ?, expires_at = ? WHERE id = ?",
407 &[
408 opt(progress.multipart_id.as_deref()),
409 text(&serde_json::to_string(&progress.parts)?),
410 num(progress.offset),
411 num(progress.tail),
412 text(&progress.hasher.save()),
413 text(&rfc3339(now_ms + DAY_MS)),
414 text(&progress.id),
415 ],
416 )?
417 .run()
418 .await?;
419 Ok(())
420 }
421
422 pub async fn delete_upload(&self, id: &str) -> Result<()> {
423 self.prepare("DELETE FROM uploads WHERE id = ?", &[text(id)])?.run().await?;
424 Ok(())
425 }
426
427 pub async fn expired_uploads(&self, now_ms: u64, limit: u32) -> Result<Vec<UploadRow>> {
428 self.prepare(
429 &format!(
430 "SELECT id, workspace, package_id, package, multipart_id, parts, \"offset\" AS offset, tail, hash_state
431 FROM uploads WHERE expires_at < ? ORDER BY expires_at LIMIT {limit}"
432 ),
433 &[text(&rfc3339(now_ms))],
434 )?
435 .all()
436 .await?
437 .results()
438 }
439
440 pub async fn version_by_digest(&self, package_id: &str, digest: &str) -> Result<Option<VersionRow>> {
441 self.prepare(
442 &format!("SELECT {VERSION_COLUMNS} FROM versions WHERE package_id = ? AND digest = ? LIMIT 1"),
443 &[text(package_id), text(digest)],
444 )?
445 .first(None)
446 .await
447 }
448
449 pub async fn version_by_tag(&self, package_id: &str, tag: &str) -> Result<Option<VersionRow>> {
450 self.prepare(
451 &format!(
452 "SELECT {} FROM tags t JOIN versions v ON v.id = t.version_id WHERE t.package_id = ? AND t.tag = ?",
453 VERSION_COLUMNS.split(", ").map(|c| format!("v.{c}")).collect::<Vec<_>>().join(", ")
454 ),
455 &[text(package_id), text(tag)],
456 )?
457 .first(None)
458 .await
459 }
460
461 /// A version by its version, digest, or a tag that points to it.
462 pub async fn find_version(&self, package_id: &str, reference: &str) -> Result<Option<VersionRow>> {
463 if let Some(found) = self.version_by_digest(package_id, reference).await? {
464 return Ok(Some(found));
465 }
466 let by_version: Option<VersionRow> = self
467 .prepare(
468 &format!("SELECT {VERSION_COLUMNS} FROM versions WHERE package_id = ? AND version = ?"),
469 &[text(package_id), text(reference)],
470 )?
471 .first(None)
472 .await?;
473 if by_version.is_some() {
474 return Ok(by_version);
475 }
476 self.version_by_tag(package_id, reference).await
477 }
478
479 pub async fn versions(&self, package_id: &str, limit: u32) -> Result<Vec<VersionRow>> {
480 self.prepare(
481 &format!("SELECT {VERSION_COLUMNS} FROM versions WHERE package_id = ? ORDER BY published_at DESC, id DESC LIMIT {limit}"),
482 &[text(package_id)],
483 )?
484 .all()
485 .await?
486 .results()
487 }
488
489 /// Records a version unless one of its digest is there, moves `tag` to
490 /// it, and says whether anything was published: a new version, or a
491 /// tag that now points somewhere else.
492 pub async fn publish(&self, version: NewVersion, tag: Option<&str>, now_ms: u64) -> Result<(VersionRow, bool)> {
493 let now = rfc3339(now_ms);
494 let existing = self.version_by_digest(&version.package_id, &version.digest).await?;
495 let mut changed = existing.is_none();
496 let mut batch = Vec::new();
497 let version_id = match &existing {
498 Some(row) => row.id.clone(),
499 None => {
500 batch.push(self.prepare(
501 "INSERT INTO versions (id, package_id, version, digest, size, metadata, subject, published_by, published_at)
502 VALUES (?, ?, ?, ?, ?, ?, ?, ?, ?)",
503 &[
504 text(&version.id),
505 text(&version.package_id),
506 text(&version.version),
507 text(&version.digest),
508 num(version.size),
509 text(&version.metadata),
510 opt(version.subject.as_deref()),
511 opt(version.published_by.as_deref()),
512 text(&now),
513 ],
514 )?);
515 for file in &version.files {
516 batch.push(self.prepare(
517 "INSERT OR IGNORE INTO version_files (version_id, name, digest, size, media_type) VALUES (?, ?, ?, ?, ?)",
518 &[text(&version.id), text(&file.name), text(&file.digest), num(file.size), opt(file.media_type.as_deref())],
519 )?);
520 }
521 version.id.clone()
522 }
523 };
524 if let Some(tag) = tag {
525 let before = self.version_by_tag(&version.package_id, tag).await?;
526 changed |= before.map(|row| row.id) != Some(version_id.clone());
527 batch.push(self.prepare(
528 "INSERT INTO tags (package_id, tag, version_id, updated_at) VALUES (?, ?, ?, ?)
529 ON CONFLICT (package_id, tag) DO UPDATE SET version_id = excluded.version_id, updated_at = excluded.updated_at",
530 &[text(&version.package_id), text(tag), text(&version_id), text(&now)],
531 )?);
532 }
533 if changed {
534 batch.push(self.prepare("UPDATE packages SET updated_at = ? WHERE id = ?", &[text(&now), text(&version.package_id)])?);
535 }
536 if !batch.is_empty() {
537 self.db.batch(batch).await?;
538 }
539 let row = self
540 .version_by_digest(&version.package_id, &version.digest)
541 .await?
542 .ok_or_else(|| worker::Error::RustError("the version was not recorded".into()))?;
543 Ok((row, changed))
544 }
545
546 pub async fn tags(&self, package_id: &str) -> Result<Vec<TagRow>> {
547 self.prepare(
548 "SELECT t.tag, t.version_id, v.digest, t.updated_at FROM tags t JOIN versions v ON v.id = t.version_id
549 WHERE t.package_id = ? ORDER BY t.tag",
550 &[text(package_id)],
551 )?
552 .all()
553 .await?
554 .results()
555 }
556
557 /// Tag names in order, `n` of them after `last`.
558 pub async fn tag_names(&self, package_id: &str, last: Option<&str>, n: u32) -> Result<Vec<String>> {
559 #[derive(Deserialize)]
560 struct Name {
561 tag: String,
562 }
563 let rows: Vec<Name> = self
564 .prepare(
565 &format!("SELECT tag FROM tags WHERE package_id = ? AND tag > ? ORDER BY tag LIMIT {n}"),
566 &[text(package_id), text(last.unwrap_or(""))],
567 )?
568 .all()
569 .await?
570 .results()?;
571 Ok(rows.into_iter().map(|row| row.tag).collect())
572 }
573
574 pub async fn referrers(&self, package_id: &str, subject: &str) -> Result<Vec<VersionRow>> {
575 self.prepare(
576 &format!("SELECT {VERSION_COLUMNS} FROM versions WHERE package_id = ? AND subject = ? ORDER BY published_at"),
577 &[text(package_id), text(subject)],
578 )?
579 .all()
580 .await?
581 .results()
582 }
583
584 pub async fn delete_tag(&self, package_id: &str, tag: &str) -> Result<()> {
585 self.prepare("DELETE FROM tags WHERE package_id = ? AND tag = ?", &[text(package_id), text(tag)])?
586 .run()
587 .await?;
588 Ok(())
589 }
590
591 /// A version, its tags and its files. Blobs are left to the sweep.
592 pub async fn delete_version(&self, version_id: &str) -> Result<()> {
593 let id = [text(version_id)];
594 self.db
595 .batch(vec![
596 self.prepare("DELETE FROM tags WHERE version_id = ?", &id)?,
597 self.prepare("DELETE FROM version_files WHERE version_id = ?", &id)?,
598 self.prepare("DELETE FROM versions WHERE id = ?", &id)?,
599 ])
600 .await?;
601 Ok(())
602 }
603
604 /// Works out again what the workspace stores: each blob its versions
605 /// use, once, public when any public package uses it.
606 pub async fn measure(&self, workspace: &str) -> Result<()> {
607 let ws = [text(workspace)];
608 self.db
609 .batch(vec![
610 self.prepare("DELETE FROM workspace_blobs WHERE workspace = ?", &ws)?,
611 self.prepare(
612 "INSERT INTO workspace_blobs (workspace, digest, size, public)
613 SELECT p.workspace, vf.digest, MAX(vf.size), MAX(p.visibility = 'public')
614 FROM version_files vf JOIN versions v ON v.id = vf.version_id JOIN packages p ON p.id = v.package_id
615 WHERE p.workspace = ? GROUP BY vf.digest",
616 &ws,
617 )?,
618 ])
619 .await?;
620 Ok(())
621 }
622
623 pub async fn storage(&self, workspace: &str) -> Result<(u64, u64)> {
624 #[derive(Deserialize)]
625 struct Sums {
626 public_bytes: Option<u64>,
627 private_bytes: Option<u64>,
628 }
629 let sums: Option<Sums> = self
630 .prepare(
631 &format!(
632 "SELECT SUM(CASE WHEN public = 1 THEN size ELSE 0 END) AS public_bytes,
633 SUM(CASE WHEN public = 1 THEN 0 ELSE size END) AS private_bytes
634 FROM workspace_blobs WHERE workspace = ? AND workspace NOT IN ({DELETED_WORKSPACES})"
635 ),
636 &[text(workspace)],
637 )?
638 .first(None)
639 .await?;
640 Ok(sums.map_or((0, 0), |s| (s.public_bytes.unwrap_or(0), s.private_bytes.unwrap_or(0))))
641 }
642
643 /// Hides (`deleted_at` given) or shows again (`None`) a workspace's
644 /// packages. Hiding keeps a mark already set; showing clears only marks.
645 pub async fn mark_workspace(&self, workspace: &str, deleted_at: Option<&str>) -> Result<()> {
646 let statement = match deleted_at {
647 Some(at) => self.prepare(
648 "UPDATE packages SET workspace_deleted_at = ? WHERE workspace = ? AND workspace_deleted_at IS NULL",
649 &[text(at), text(workspace)],
650 )?,
651 None => self.prepare(
652 "UPDATE packages SET workspace_deleted_at = NULL WHERE workspace = ? AND workspace_deleted_at IS NOT NULL",
653 &[text(workspace)],
654 )?,
655 };
656 statement.run().await?;
657 Ok(())
658 }
659
660 /// Whether the workspace is deleted, as its packages say.
661 pub async fn workspace_hidden(&self, workspace: &str) -> Result<bool> {
662 let row: Option<serde_json::Value> = self
663 .prepare("SELECT 1 AS hidden FROM packages WHERE workspace = ? AND workspace_deleted_at IS NOT NULL LIMIT 1", &[text(workspace)])?
664 .first(None)
665 .await?;
666 Ok(row.is_some())
667 }
668
669 /// Every workspace's package storage, by workspace.
670 pub async fn storage_all(&self) -> Result<Vec<g1t_contracts::packages::WorkspacePackageStorage>> {
671 self.prepare(
672 &format!(
673 "SELECT workspace,
674 COALESCE(SUM(CASE WHEN public = 1 THEN size ELSE 0 END), 0) AS public_bytes,
675 COALESCE(SUM(CASE WHEN public = 1 THEN 0 ELSE size END), 0) AS private_bytes
676 FROM workspace_blobs WHERE workspace NOT IN ({DELETED_WORKSPACES}) GROUP BY workspace ORDER BY workspace"
677 ),
678 &[],
679 )?
680 .all()
681 .await?
682 .results()
683 }
684
685 /// Which of `digests` the workspace already holds.
686 pub async fn held(&self, workspace: &str, digests: &[String]) -> Result<std::collections::HashSet<String>> {
687 #[derive(Deserialize)]
688 struct Held {
689 digest: String,
690 }
691 let mut held = std::collections::HashSet::new();
692 for chunk in digests.chunks(50) {
693 let marks = vec!["?"; chunk.len()].join(", ");
694 let mut values = vec![text(workspace)];
695 values.extend(chunk.iter().map(|d| text(d)));
696 let rows: Vec<Held> = self
697 .prepare(&format!("SELECT digest FROM workspace_blobs WHERE workspace = ? AND digest IN ({marks})"), &values)?
698 .all()
699 .await?
700 .results()?;
701 held.extend(rows.into_iter().map(|row| row.digest));
702 }
703 Ok(held)
704 }
705
706 /// Blobs no version uses and nothing has touched for a day.
707 pub async fn unused_blobs(&self, now_ms: u64, limit: u32) -> Result<Vec<BlobRow>> {
708 self.prepare(
709 &format!(
710 "SELECT digest, size, media_type, object_key FROM blobs b
711 WHERE touched_at < ? AND NOT EXISTS (SELECT 1 FROM version_files vf WHERE vf.digest = b.digest)
712 ORDER BY touched_at LIMIT {limit}"
713 ),
714 &[text(&rfc3339(now_ms.saturating_sub(DAY_MS)))],
715 )?
716 .all()
717 .await?
718 .results()
719 }
720
721 pub async fn forget_blob(&self, digest: &str) -> Result<()> {
722 let d = [text(digest)];
723 self.db
724 .batch(vec![
725 self.prepare("DELETE FROM package_blobs WHERE digest = ?", &d)?,
726 self.prepare("DELETE FROM workspace_blobs WHERE digest = ?", &d)?,
727 self.prepare("DELETE FROM blobs WHERE digest = ?", &d)?,
728 ])
729 .await?;
730 Ok(())
731 }
732
733 /// Workspaces with packages linked to the repository.
734 pub async fn workspaces_linked_to(&self, repo_id: &str) -> Result<Vec<String>> {
735 #[derive(Deserialize)]
736 struct Ws {
737 workspace: String,
738 }
739 let rows: Vec<Ws> = self
740 .prepare("SELECT DISTINCT workspace FROM packages WHERE repo_id = ?", &[text(repo_id)])?
741 .all()
742 .await?
743 .results()?;
744 Ok(rows.into_iter().map(|row| row.workspace).collect())
745 }
746
747 pub async fn follow_visibility(&self, repo_id: &str, private: bool) -> Result<()> {
748 self.prepare(
749 "UPDATE packages SET visibility = ? WHERE repo_id = ?",
750 &[text(if private { "private" } else { "public" }), text(repo_id)],
751 )?
752 .run()
753 .await?;
754 Ok(())
755 }
756}