g1t/services/packages/src/db.rs

930 lines36,560 bytesCodeBlame

Pick any line to see why it is the way it is: the commit, the pull request and issue it came from, and what the agent was thinking.

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