g1t/services/packages/src/db.rs

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