Skip to content

g1t/services/packages/src/lib.rs

1,057 lines46,978 bytesCodeBlameRaw
1//! The packages service: the registries a workspace publishes to and
2//! installs from, beside its code (docs/PACKAGES.md). Container images
3//! first, over OCI Distribution 1.1 on `g1t.sh/v2/` (oci.rs).
4//!
5//! Other services reach it over `POST /rpc/<method>`; see
6//! `g1t_contracts::packages` for the methods and their arguments. Any
7//! other request is a registry's own protocol. Files are kept through the
8//! `BlobStore` port (store/), metadata in D1 (db.rs).
9
10mod access;
11mod archive;
12mod cargo;
13mod cargo_http;
14mod composer;
15mod composer_http;
16mod db;
17mod digest;
18mod limits;
19mod manifest;
20mod marshal;
21mod maven;
22mod maven_http;
23mod names;
24mod npm;
25mod npm_http;
26mod nuget;
27mod nuget_http;
28mod oci;
29mod quota;
30mod range;
31mod rubygems;
32mod rubygems_http;
33mod store;
34mod token;
35mod upload;
36mod xml;
37mod yaml;
38
39use std::cell::RefCell;
40use std::collections::HashMap;
41
42use g1t_contracts::audit::{AuditActor, AuditTarget, NewAuditEntry, RecordAuditArgs, Surface};
43use g1t_contracts::credentials::Decision;
44use g1t_contracts::events::{Event, NewEvent, PackageEvent, Publish};
45use g1t_contracts::identity::GitCredentialsArgs;
46use g1t_contracts::packages::*;
47use g1t_contracts::repos::{GetArgs, Repo, RepoPath};
48use g1t_contracts::{FailureCode, Outcome, User, Viewer, new_id};
49use g1t_kit::{args, now_ms, reply, rpc_method};
50use worker::{Context, Env, Fetcher, MessageBatch, Request, Response, Result, ScheduleContext, ScheduledEvent, event};
51
52use access::{Action, LinkedTo, Target};
53use db::{Db, PackageRow};
54use store::{BlobStore, Store};
55
56mod settings;
57
58pub(crate) const SOURCE: &str = "packages";
59/// How large a request body may be: Cloudflare's limit on the zone's plan.
60const DEFAULT_MAX_REQUEST_BYTES: u64 = 100_000_000;
61/// How many packages a listing shows.
62const LIST_LIMIT: u32 = 200;
63const VERSIONS_SHOWN: u32 = 200;
64/// How much one sweep lets go of.
65const SWEEP_BATCH: u32 = 200;
66
67thread_local! {
68 /// Pulls counted since the last write, by package (and by version,
69 /// where it is counted too): written at most every few seconds, so a
70 /// busy image costs one write, not one a pull. What an isolate holds
71 /// when it goes away is lost: the count is approximate.
72 static DOWNLOADS: RefCell<(HashMap<DownloadKey, u64>, u64)> = RefCell::new((HashMap::new(), 0));
73}
74const DOWNLOADS_FLUSH_MS: u64 = 10_000;
75/// A package's id, and a version's when the download counts for it too.
76type DownloadKey = (String, Option<String>);
77
78thread_local! {
79 /// What billing allows each workspace, as asked last, and when.
80 static ALLOWANCES: RefCell<HashMap<String, (quota::Allowance, u64)>> = RefCell::new(HashMap::new());
81}
82/// How long billing's answer is kept.
83const ALLOWANCE_TTL_MS: u64 = 5 * 60 * 1000;
84
85/// Who made a registry request, as its audit entries and versions name them.
86pub struct Caller {
87 pub actor: Option<AuditActor>,
88 /// Who it is, where that is known: what linking a package to the
89 /// repository its source label names checks.
90 pub viewer: Option<User>,
91}
92
93impl Caller {
94 pub fn of(viewer: Option<&User>) -> Caller {
95 Caller { actor: viewer.map(AuditActor::of), viewer: viewer.cloned() }
96 }
97
98 pub fn system() -> Caller {
99 Caller { actor: Some(AuditActor::system()), viewer: None }
100 }
101}
102
103/// What access decisions need about a package, owned. Its own grants and
104/// Actions access are read only when a decision needs them (`loaded`).
105pub struct TargetOf {
106 pub(crate) workspace: String,
107 pub(crate) name: String,
108 pub(crate) repo: Option<(String, String, bool)>,
109 pub(crate) public: bool,
110 pub(crate) package_id: Option<String>,
111 pub(crate) inherit: bool,
112 pub(crate) grants: Vec<access::Grant>,
113 pub(crate) actions: Vec<access::RepoAccess>,
114 pub(crate) teams: Vec<String>,
115 pub(crate) loaded: bool,
116}
117
118impl TargetOf {
119 pub fn package(row: &PackageRow) -> TargetOf {
120 TargetOf {
121 workspace: row.workspace.clone(),
122 name: row.name.clone(),
123 repo: row
124 .repo_id
125 .as_ref()
126 .map(|id| (id.clone(), row.repo_name.clone().unwrap_or_default(), !row.public())),
127 public: row.public(),
128 package_id: Some(row.id.clone()),
129 inherit: row.inherits(),
130 grants: Vec::new(),
131 actions: Vec::new(),
132 teams: Vec::new(),
133 loaded: false,
134 }
135 }
136
137 /// A package not made yet, which its first push would link to `repo`.
138 pub fn unmade(workspace: &str, name: &str, repo: Option<(String, String, bool)>) -> TargetOf {
139 TargetOf {
140 workspace: workspace.to_owned(),
141 name: name.to_owned(),
142 repo,
143 public: false,
144 package_id: None,
145 inherit: true,
146 grants: Vec::new(),
147 actions: Vec::new(),
148 teams: Vec::new(),
149 loaded: true,
150 }
151 }
152
153 pub fn view(&self) -> Target<'_> {
154 Target {
155 workspace: &self.workspace,
156 name: &self.name,
157 repo: self.repo.as_ref().map(|(id, name, private)| LinkedTo { id, name, private: *private }),
158 public: self.public,
159 exists: self.package_id.is_some(),
160 inherit: self.inherit,
161 grants: &self.grants,
162 actions: &self.actions,
163 teams: &self.teams,
164 }
165 }
166}
167
168pub struct Packages {
169 pub db: Db,
170 pub store: Store,
171 identity: Fetcher,
172 repos: Fetcher,
173 events: Fetcher,
174 /// Signs registry tokens (PACKAGES_TOKEN_SECRET).
175 secret: Vec<u8>,
176 max_request: u64,
177 /// The host package addresses start with: `g1t.sh`.
178 host: String,
179 /// Whether free workspaces' package storage is limited (STORAGE_LIMITS
180 /// `on`); off when self-hosted.
181 storage_limits: bool,
182 env: Env,
183}
184
185/// What a workspace's deletion or restoring does to its packages.
186#[derive(Debug, PartialEq, Eq)]
187pub(crate) enum WorkspaceMark {
188 /// Hide them from `at`, while the workspace may still be restored.
189 Hide { workspace: String, at: String },
190 /// Show them again.
191 Show { workspace: String },
192}
193
194/// The mark `event` asks for, if any. A protected workspace (`protected`,
195/// slugs or ids, as identity's PROTECTED_WORKSPACES) is never hidden: it
196/// cannot be deleted, so an event saying so is a mistake.
197pub(crate) fn workspace_mark(event: &Event, protected: &[String], now: &str) -> Option<WorkspaceMark> {
198 let text = |key: &str| event.data[key].as_str().map(|v| v.trim().to_lowercase()).filter(|v| !v.is_empty());
199 let workspace = text("slug")?;
200 match event.kind.as_str() {
201 "workspace.deleting" => {
202 let id = text("workspaceId").unwrap_or_default();
203 if protected.iter().any(|name| *name == workspace || (!id.is_empty() && *name == id)) {
204 return None;
205 }
206 Some(WorkspaceMark::Hide { workspace, at: now.to_owned() })
207 }
208 "workspace.restored" => Some(WorkspaceMark::Show { workspace }),
209 _ => None,
210 }
211}
212
213/// Adds one download to the counts kept since the last write, and answers
214/// with all of them to write when the last write was long enough ago.
215pub(crate) fn tally(counts: &mut (HashMap<DownloadKey, u64>, u64), key: DownloadKey, now: u64) -> Option<Vec<(DownloadKey, u64)>> {
216 *counts.0.entry(key).or_default() += 1;
217 if now.saturating_sub(counts.1) < DOWNLOADS_FLUSH_MS {
218 return None;
219 }
220 counts.1 = now;
221 Some(counts.0.drain().collect())
222}
223
224fn not_found<T>() -> Outcome<T> {
225 Outcome::fail(FailureCode::NotFound, "Package not found.")
226}
227
228impl Packages {
229 pub fn from_env(env: &Env) -> Result<Packages> {
230 let secret = store::var(env, "PACKAGES_TOKEN_SECRET");
231 if secret.len() < 32 {
232 return Err(worker::Error::RustError("PACKAGES_TOKEN_SECRET is missing or shorter than 32 characters".into()));
233 }
234 // 0 is no limit, as self-hosted (nothing in front cuts bodies short).
235 let max_request = match store::var(env, "MAX_REQUEST_BYTES").parse::<u64>() {
236 Ok(0) => u64::MAX,
237 Ok(limit) => limit,
238 Err(_) => DEFAULT_MAX_REQUEST_BYTES,
239 };
240 let host = store::var(env, "REGISTRY_HOST");
241 Ok(Packages {
242 db: Db { db: env.d1("DB")? },
243 store: store::from_env(env)?,
244 identity: env.service("IDENTITY")?,
245 repos: env.service("REPOS")?,
246 events: env.service("EVENTS")?,
247 secret: secret.into_bytes(),
248 max_request,
249 host: if host.is_empty() { "g1t.sh".to_owned() } else { host },
250 storage_limits: store::var(env, "STORAGE_LIMITS") == "on",
251 env: env.clone(),
252 })
253 }
254
255 async fn viewer_for(&self, username: &str, secret: &str) -> Result<Viewer> {
256 g1t_kit::call(
257 &self.identity,
258 "user_for_git_credentials",
259 &GitCredentialsArgs { username: username.to_owned(), secret: secret.to_owned() },
260 )
261 .await
262 }
263
264 /// A repository of the workspace by name, whoever may see it.
265 async fn repo_by_name(&self, workspace: &str, name: &str) -> Result<Option<Repo>> {
266 if !g1t_contracts::is_valid_repo_name(name) {
267 return Ok(None);
268 }
269 let found: Outcome<Repo> = g1t_kit::call(
270 &self.repos,
271 "get",
272 &GetArgs {
273 path: RepoPath { namespace: workspace.to_owned(), name: name.to_owned() },
274 viewer: Some(User::system(workspace)),
275 },
276 )
277 .await?;
278 Ok(found.into_result().ok().filter(|repo| repo.fork_of.is_none()))
279 }
280
281 /// Who may do what to an image: the package's own settings, or, before
282 /// its first push, the repository it will be linked to.
283 async fn target(&self, name: &names::ImageName, found: Option<&PackageRow>) -> Result<TargetOf> {
284 if let Some(found) = found {
285 return Ok(TargetOf::package(found));
286 }
287 let repo = self.repo_by_name(&name.workspace, name.repo_name()).await?;
288 Ok(TargetOf::unmade(&name.workspace, &name.name, repo.map(|repo| (repo.id, repo.name, repo.is_private))))
289 }
290
291 /// What billing allows the workspace, kept for five minutes unless
292 /// `fresh`, and whether it is the kept answer. `None` when billing
293 /// cannot be asked: the push is then let through.
294 async fn allowance(&self, workspace: &str, fresh: bool) -> Option<(quota::Allowance, bool)> {
295 let now = now_ms();
296 let kept = ALLOWANCES.with(|kept| kept.borrow().get(workspace).copied());
297 if let Some((allowance, at)) = kept
298 && !fresh
299 && now.saturating_sub(at) < ALLOWANCE_TTL_MS
300 {
301 return Some((allowance, true));
302 }
303 let billing = self.env.service("BILLING").ok()?;
304 let asked: Result<quota::Allowance> = g1t_kit::call(
305 &billing,
306 "entitlements",
307 &g1t_contracts::billing::EntitlementsArgs { workspace: workspace.to_owned() },
308 )
309 .await;
310 match asked {
311 Ok(allowance) => {
312 ALLOWANCES.with(|kept| kept.borrow_mut().insert(workspace.to_owned(), (allowance, now)));
313 Some((allowance, false))
314 }
315 Err(error) => {
316 worker::console_error!("packages: billing could not be asked about {workspace}, letting the push through: {error}");
317 None
318 }
319 }
320 }
321
322 /// Why a push of these files (digest and size) into the package may
323 /// not be kept, if it may not: a free workspace past its free storage.
324 /// Files the workspace already holds add nothing.
325 pub(crate) async fn storage_refusal(&self, package: &PackageRow, files: &[(String, u64)]) -> Result<Option<String>> {
326 if !self.storage_limits || files.is_empty() {
327 return Ok(None);
328 }
329 let digests: Vec<String> = files.iter().map(|(digest, _)| digest.clone()).collect();
330 let held = self.db.held(&package.workspace, &digests).await?;
331 let mut seen = std::collections::HashSet::new();
332 let adding: u64 = files
333 .iter()
334 .filter(|(digest, _)| !held.contains(digest) && seen.insert(digest.clone()))
335 .map(|(_, size)| size)
336 .sum();
337 if adding == 0 {
338 return Ok(None);
339 }
340 let Some((allowance, kept)) = self.allowance(&package.workspace, false).await else {
341 return Ok(None);
342 };
343 let (public_bytes, private_bytes) = self.db.storage(&package.workspace).await?;
344 let public = package.public();
345 let used = if public { public_bytes } else { private_bytes };
346 let mut refused = quota::decide(&allowance, public, used, adding).err();
347 // A refusal from the kept answer is checked with billing again: the
348 // workspace may have just added a plan, and must not wait minutes
349 // for the push to go through.
350 if refused.is_some() && kept {
351 refused = match self.allowance(&package.workspace, true).await {
352 Some((allowance, _)) => quota::decide(&allowance, public, used, adding).err(),
353 None => None,
354 };
355 }
356 Ok(refused.map(|refusal| quota::message(&package.workspace, &refusal)))
357 }
358
359 fn count_download(&self, package_id: &str, ctx: &Context) {
360 self.count_downloads(package_id, None, ctx);
361 }
362
363 /// A download of one version, counted for it and its package.
364 fn count_version_download(&self, package_id: &str, version_id: &str, ctx: &Context) {
365 self.count_downloads(package_id, Some(version_id), ctx);
366 }
367
368 fn count_downloads(&self, package_id: &str, version_id: Option<&str>, ctx: &Context) {
369 let key = (package_id.to_owned(), version_id.map(str::to_owned));
370 let due = DOWNLOADS.with(|counts| tally(&mut counts.borrow_mut(), key, now_ms()));
371 if let Some(due) = due {
372 let env = self.env.clone();
373 ctx.wait_until(async move {
374 let written = match env.d1("DB") {
375 Ok(db) => Db { db }.add_downloads(&due).await,
376 Err(error) => Err(error),
377 };
378 if let Err(error) = written {
379 worker::console_error!("packages: downloads not counted: {error}");
380 }
381 });
382 }
383 }
384
385 async fn announce(&self, kind: &'static str, package: &PackageRow, data: PackageEvent, caller: &Caller) {
386 let event = NewEvent {
387 kind,
388 source: SOURCE,
389 repo_id: package.repo_id.clone(),
390 actor: caller.actor.as_ref().map(|actor| actor.actor_id.clone()),
391 data,
392 };
393 let published: Result<serde_json::Value> = g1t_kit::call(&self.events, "publish", &Publish { events: vec![event] }).await;
394 if let Err(error) = published {
395 worker::console_error!("packages: {kind} not published: {error}");
396 }
397 }
398
399 async fn audit(&self, caller: &Caller, action: &str, package: &PackageRow, path: Option<&str>, surface: Option<Surface>) {
400 self.audit_with(caller, action, package, path, surface, None).await;
401 }
402
403 /// An audit entry, with what changed in words (`message`), as identity
404 /// writes its access entries.
405 async fn audit_with(
406 &self,
407 caller: &Caller,
408 action: &str,
409 package: &PackageRow,
410 path: Option<&str>,
411 surface: Option<Surface>,
412 message: Option<String>,
413 ) {
414 let Some(actor) = caller.actor.clone() else {
415 return;
416 };
417 let target = AuditTarget {
418 workspace: package.workspace.clone(),
419 repo: package.repo_name.as_ref().map(|name| format!("{}/{name}", package.workspace)),
420 path: Some(path.map_or_else(|| format!("{}:{}/{}", package.ecosystem, package.workspace, package.name), str::to_owned)),
421 ..AuditTarget::default()
422 };
423 let mut entry = NewAuditEntry::new(
424 actor,
425 action,
426 surface.unwrap_or(Surface::Registry),
427 target,
428 &Decision::allow("packages"),
429 new_id("req", now_ms()),
430 );
431 entry.result = Some("ok".to_owned());
432 entry.message = message;
433 let recorded: Result<u32> = g1t_kit::call(&self.events, "audit_record", &RecordAuditArgs { entries: vec![entry] }).await;
434 if let Err(error) = recorded {
435 worker::console_error!("packages: audit entry not recorded: {error}");
436 }
437 }
438
439 fn summary(&self, row: &db::ListedRow) -> PackageSummary {
440 let p = &row.package;
441 PackageSummary {
442 id: p.id.clone(),
443 workspace: p.workspace.clone(),
444 ecosystem: Ecosystem::parse(&p.ecosystem).unwrap_or(Ecosystem::Container),
445 name: p.name.clone(),
446 address: match p.ecosystem.as_str() {
447 "npm" => format!("{}/-/npm/@{}/{}", self.host, p.workspace, p.name),
448 "composer" => format!("{}/-/composer/{}/{}", self.host, p.workspace, p.name),
449 "cargo" => format!("{}/-/cargo/{}/{}", self.host, p.workspace, p.name),
450 "maven" => format!("{}/-/maven/{}/{}", self.host, p.workspace, p.name),
451 "nuget" => format!("{}/-/nuget/{}/{}", self.host, p.workspace, p.name),
452 "rubygems" => format!("{}/-/rubygems/{}/{}", self.host, p.workspace, p.name),
453 _ => format!("{}/{}/{}", self.host, p.workspace, p.name),
454 },
455 visibility: Visibility::parse(&p.visibility),
456 repo: p.repo_id.as_ref().map(|id| LinkedRepo {
457 id: id.clone(),
458 namespace: p.workspace.clone(),
459 name: p.repo_name.clone().unwrap_or_default(),
460 }),
461 description: p.description.clone(),
462 versions: row.version_count,
463 latest: db::latest_shown(row),
464 size: row.bytes,
465 downloads: p.downloads,
466 created_at: p.created_at.clone(),
467 updated_at: p.updated_at.clone(),
468 inherit_access: p.inherits(),
469 deleted_at: p.deleted_at.clone(),
470 deleted_by: p.deleted_by.clone(),
471 purge_at: p.deleted_at.as_deref().and_then(settings::purge_at),
472 }
473 }
474
475 async fn listed(&self, workspace: &str, ecosystem: Ecosystem, name: &str) -> Result<Option<db::ListedRow>> {
476 let rows = self.db.list(workspace, Some(ecosystem.as_str()), None, None, LIST_LIMIT).await?;
477 if let Some(row) = rows.into_iter().find(|row| row.package.name == name) {
478 return Ok(Some(row));
479 }
480 // Past the first page: read it alone.
481 Ok(self.db.package(workspace, ecosystem.as_str(), name).await?.filter(|package| !package.hidden()).map(|package| db::ListedRow {
482 package,
483 version_count: 0,
484 bytes: 0,
485 latest_tag: None,
486 latest_tag_version: None,
487 latest_version: None,
488 }))
489 }
490
491 async fn list_packages(&self, a: ListPackagesArgs) -> Result<Outcome<Vec<PackageSummary>>> {
492 let workspace = a.workspace.to_lowercase();
493 let rows = self
494 .db
495 .list(&workspace, a.ecosystem.map(Ecosystem::as_str), a.repo_id.as_deref(), a.query.as_deref(), LIST_LIMIT)
496 .await?;
497 let packages: Vec<&PackageRow> = rows.iter().map(|row| &row.package).collect();
498 let may = self.may_all(a.viewer.as_ref(), &packages, Action::Pull).await?;
499 let visible = rows
500 .iter()
501 .zip(may)
502 .filter(|(_, may)| *may)
503 .map(|(row, _)| self.summary(row))
504 .collect();
505 Ok(Outcome::Ok(visible))
506 }
507
508 async fn get_package(&self, a: GetPackageArgs) -> Result<Outcome<PackageDetail>> {
509 let workspace = a.workspace.to_lowercase();
510 let Some(row) = self.listed(&workspace, a.ecosystem, &a.name).await? else {
511 return Ok(not_found());
512 };
513 let mut target = TargetOf::package(&row.package);
514 let permissions = self.permissions(a.viewer.as_ref(), &mut target).await?;
515 if !permissions.pull {
516 return Ok(not_found());
517 }
518 let tags = self.db.tags(&row.package.id).await?;
519 let versions = self
520 .db
521 .versions(&row.package.id, VERSIONS_SHOWN)
522 .await?
523 .into_iter()
524 .map(|version| settings::version_of(&row.package, version, &tags))
525 .collect();
526 // The README its page shows: npm's, from the latest version.
527 let readme = match self.db.readme_digest(&row.package.id).await?.and_then(|d| digest::Digest::parse(&d)) {
528 Some(digest) => match self.db.blob(&digest).await? {
529 Some(blob) => self.store.read(&blob.object_key).await?.map(|b| String::from_utf8_lossy(&b).into_owned()),
530 None => None,
531 },
532 None => None,
533 };
534 Ok(Outcome::Ok(PackageDetail {
535 readme,
536 package: self.summary(&row),
537 versions,
538 tags: tags
539 .into_iter()
540 .map(|tag| PackageTag { tag: tag.tag, digest: tag.digest, updated_at: tag.updated_at })
541 .collect(),
542 permissions,
543 }))
544 }
545
546 /// The package `actor` asks to change, if they may `action` it.
547 pub(crate) async fn for_change(&self, actor: &User, workspace: &str, ecosystem: Ecosystem, name: &str, action: Action) -> Result<Outcome<PackageRow>> {
548 let Some(package) = self.db.package(&workspace.to_lowercase(), ecosystem.as_str(), name).await?.filter(|p| !p.hidden()) else {
549 return Ok(not_found());
550 };
551 match self.allowed(actor, &package, action).await? {
552 Outcome::Ok(()) => Ok(Outcome::Ok(package)),
553 Outcome::Fail(failure) => Ok(Outcome::Fail(failure)),
554 }
555 }
556
557 /// Whether `actor` may `action` the package: not found for one who may
558 /// not pull it, forbidden with the reason for one who may.
559 pub(crate) async fn allowed(&self, actor: &User, package: &PackageRow, action: Action) -> Result<Outcome<()>> {
560 let mut target = TargetOf::package(package);
561 let decision = self.decide(Some(actor), &mut target, action).await?;
562 if decision.allowed {
563 return Ok(Outcome::Ok(()));
564 }
565 if !self.decide(Some(actor), &mut target, Action::Pull).await?.allowed {
566 return Ok(not_found());
567 }
568 Ok(Outcome::fail(FailureCode::Forbidden, decision.reason.unwrap_or_else(|| "Not allowed.".to_owned())))
569 }
570
571 async fn delete_version(&self, a: DeleteVersionArgs) -> Result<Outcome<()>> {
572 let package = match self.for_change(&a.actor, &a.workspace, a.ecosystem, &a.name, Action::Delete).await? {
573 Outcome::Ok(package) => package,
574 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
575 };
576 if package.ecosystem == "composer" {
577 return Ok(Outcome::fail(
578 FailureCode::Invalid,
579 "A Composer package's versions are its repository's tags and branches: delete the tag or branch instead.",
580 ));
581 }
582 let Some(version) = self.db.find_version(&package.id, &a.version).await? else {
583 return Ok(Outcome::fail(FailureCode::NotFound, "Version not found."));
584 };
585 let caller = Caller::of(Some(&a.actor));
586 self.remove_version_from(&package, &version, &caller, a.surface).await?;
587 Ok(Outcome::Ok(()))
588 }
589
590 async fn delete_package(&self, a: DeletePackageArgs) -> Result<Outcome<()>> {
591 let package = match self.for_change(&a.actor, &a.workspace, a.ecosystem, &a.name, Action::Delete).await? {
592 Outcome::Ok(package) => package,
593 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
594 };
595 let caller = Caller::of(Some(&a.actor));
596 self.remove_package(&package, &caller, a.surface).await?;
597 Ok(Outcome::Ok(()))
598 }
599
600 /// Deletes a package: hidden at once, restorable for 30 days, its name
601 /// kept until then. Its files go with the purge.
602 pub(crate) async fn remove_package(&self, package: &PackageRow, caller: &Caller, surface: Option<Surface>) -> Result<()> {
603 let by = caller.actor.as_ref().map_or("", |actor| actor.actor.as_str());
604 self.db.soft_delete_package(&package.id, by, now_ms()).await?;
605 self.db.measure(&package.workspace).await?;
606 self.announce("package.deleted", package, self.event_of(package), caller).await;
607 self.audit_with(
608 caller,
609 "package.delete",
610 package,
611 None,
612 surface,
613 Some(format!("Deleted the package; it can be restored for {RESTORE_DAYS} days")),
614 )
615 .await;
616 Ok(())
617 }
618
619 async fn set_package(&self, a: SetPackageArgs) -> Result<Outcome<PackageSummary>> {
620 let package = match self.for_change(&a.actor, &a.workspace, a.ecosystem, &a.name, Action::Admin).await? {
621 Outcome::Ok(package) => package,
622 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
623 };
624 let caller = Caller::of(Some(&a.actor));
625 let now = now_ms();
626 let before = package.visibility.clone();
627 let linked = if a.unlink { false } else { a.link.is_some() || package.repo_id.is_some() };
628 if a.visibility.is_some_and(|visibility| visibility.as_str() != package.visibility) && linked {
629 return Ok(Outcome::fail(
630 FailureCode::Invalid,
631 "A linked package has its repository's visibility. Change the repository's, or unlink the package.",
632 ));
633 }
634 if a.inherit_access.is_some() && !linked {
635 return Ok(Outcome::fail(
636 FailureCode::Invalid,
637 "Only a package linked to a repository inherits access. Link it to a repository first.",
638 ));
639 }
640 let path = |repo: &str| format!("{}/{repo}", package.workspace);
641 if a.unlink {
642 if let Some(name) = package.repo_name.as_deref() {
643 self.db.set_link(&package.id, None, "private", now).await?;
644 self.audit_with(&caller, "package.unlinked", &package, None, a.surface, Some(format!("Unlinked from {}", path(name)))).await;
645 }
646 } else if let Some(link) = a.link.as_deref() {
647 let link = link.trim().trim_end_matches(".git");
648 let name = match link.split_once('/') {
649 Some((owner, name)) if owner.eq_ignore_ascii_case(&package.workspace) => name,
650 Some(_) => {
651 return Ok(Outcome::fail(
652 FailureCode::Invalid,
653 format!("A package can only be linked to a repository of its own workspace, {}.", package.workspace),
654 ));
655 }
656 None => link,
657 };
658 let Some(repo) = self.repo_by_name(&package.workspace, name).await? else {
659 return Ok(Outcome::fail(FailureCode::NotFound, format!("There is no repository {}.", path(name))));
660 };
661 // Linking hands the package to the repository's roles: only
662 // someone who administers that repository may.
663 let role = g1t_contracts::access::permission(Some(&a.actor), &repo);
664 if role < Some(g1t_contracts::access::RepoRole::Admin) {
665 return Ok(Outcome::fail(
666 FailureCode::Forbidden,
667 format!("You need the Admin role on {}/{} to link a package to it.", repo.namespace, repo.name),
668 ));
669 }
670 if package.repo_id.as_deref() != Some(repo.id.as_str()) {
671 let visibility = if repo.is_private { "private" } else { "public" };
672 self.db.set_link(&package.id, Some((&repo.id, &repo.name)), visibility, now).await?;
673 self.audit_with(&caller, "package.linked", &package, None, a.surface, Some(format!("Linked to {}", path(&repo.name)))).await;
674 }
675 }
676 if let Some(visibility) = a.visibility
677 && !linked
678 && visibility.as_str() != package.visibility
679 {
680 self.db.set_visibility(&package.id, visibility.as_str(), now).await?;
681 }
682 if let Some(inherit) = a.inherit_access
683 && inherit != package.inherits()
684 {
685 self.db.set_inherit(&package.id, inherit, now).await?;
686 let message = if inherit {
687 "Turned on inheriting access from the linked repository"
688 } else {
689 "Turned off inheriting access from the linked repository"
690 };
691 self.audit_with(&caller, "package.inherit_access_changed", &package, None, a.surface, Some(message.to_owned())).await;
692 }
693 let Some(row) = self.listed(&package.workspace, a.ecosystem, &package.name).await? else {
694 return Ok(not_found());
695 };
696 if row.package.visibility != before {
697 self.db.measure(&package.workspace).await?;
698 let event = PackageEvent { visibility: Some(row.package.visibility.clone()), ..self.event_of(&row.package) };
699 self.announce("package.visibility_changed", &row.package, event, &caller).await;
700 self.audit_with(
701 &caller,
702 "package.visibility_changed",
703 &row.package,
704 None,
705 a.surface,
706 Some(format!("Made the package {}", row.package.visibility)),
707 )
708 .await;
709 }
710 Ok(Outcome::Ok(self.summary(&row)))
711 }
712
713 /// Lets go of expired uploads, and of blobs no version has used for a day.
714 async fn sweep(&self) -> Result<(u32, u32)> {
715 let now = now_ms();
716 let mut uploads = 0;
717 for row in self.db.expired_uploads(now, SWEEP_BATCH).await? {
718 if let Some(progress) = row.progress()
719 && let Err(error) = upload::abort(&self.store, &progress).await
720 {
721 worker::console_error!("packages: upload {} not aborted: {error}", row.id);
722 }
723 self.db.delete_upload(&row.id).await?;
724 uploads += 1;
725 }
726 let mut blobs = 0;
727 for blob in self.db.unused_blobs(now, SWEEP_BATCH).await? {
728 self.store.delete(&blob.object_key).await?;
729 self.db.forget_blob(&blob.digest).await?;
730 blobs += 1;
731 }
732 Ok((uploads, blobs))
733 }
734
735 /// Follows what happens elsewhere: workspaces renamed and deleted,
736 /// repositories that change visibility, are renamed, move or go.
737 async fn on_event(&self, env: &Env, event: &Event) -> Result<()> {
738 // Composer packages are read from repositories: again when one is
739 // pushed to (its default branch may have gained a composer.json),
740 // restored, renamed or moved; gone when it is deleted. A failure is
741 // logged, not retried with the batch: the next push reads it again.
742 let repo_id = event.repo_id.clone().or_else(|| event.data["repoId"].as_str().map(str::to_owned));
743 if let Some(repo_id) = repo_id.as_deref() {
744 let synced = match event.kind.as_str() {
745 "git.push" => {
746 let known = self.db.package_for_repo(repo_id, "composer").await?.is_some();
747 if known || event.data["defaultBranch"].as_bool() == Some(true) {
748 Some(self.sync_composer(repo_id).await.map(|_| ()))
749 } else {
750 None
751 }
752 }
753 "repo.deleted" => Some(self.composer_repo_gone(repo_id).await),
754 "repo.restored" | "repo.renamed" | "repo.transferred" => Some(self.sync_composer(repo_id).await.map(|_| ())),
755 _ => None,
756 };
757 if let Some(Err(error)) = synced {
758 worker::console_error!("packages: composer {} for {repo_id}: {error}", event.kind);
759 }
760 if event.kind == "git.push" || event.kind == "repo.deleted" || event.kind == "repo.restored" {
761 return Ok(());
762 }
763 }
764 let db = &self.db.db;
765 let protected = g1t_contracts::identity::protected_names(Some(&store::var(env, "PROTECTED_WORKSPACES")));
766 match workspace_mark(event, &protected, &g1t_contracts::time::rfc3339(now_ms())) {
767 Some(WorkspaceMark::Hide { workspace, at }) => return self.db.mark_workspace(&workspace, Some(&at)).await,
768 Some(WorkspaceMark::Show { workspace }) => return self.db.mark_workspace(&workspace, None).await,
769 None => {}
770 }
771 if g1t_kit::rename::on_event(env, db, event, &[
772 "UPDATE packages SET workspace = ?1 WHERE workspace = ?2",
773 "UPDATE uploads SET workspace = ?1 WHERE workspace = ?2",
774 "UPDATE workspace_blobs SET workspace = ?1 WHERE workspace = ?2",
775 ])
776 .await?
777 {
778 return Ok(());
779 }
780 if g1t_kit::deleted::on_event(db, event, &[
781 "DELETE FROM tags WHERE package_id IN (SELECT id FROM packages WHERE workspace = ?1)",
782 "DELETE FROM version_files WHERE version_id IN (SELECT v.id FROM versions v JOIN packages p ON p.id = v.package_id WHERE p.workspace = ?1)",
783 "DELETE FROM versions WHERE package_id IN (SELECT id FROM packages WHERE workspace = ?1)",
784 "DELETE FROM package_blobs WHERE package_id IN (SELECT id FROM packages WHERE workspace = ?1)",
785 "DELETE FROM package_access WHERE package_id IN (SELECT id FROM packages WHERE workspace = ?1)",
786 "DELETE FROM package_actions_access WHERE package_id IN (SELECT id FROM packages WHERE workspace = ?1)",
787 "DELETE FROM uploads WHERE workspace = ?1",
788 "DELETE FROM workspace_blobs WHERE workspace = ?1",
789 "DELETE FROM packages WHERE workspace = ?1",
790 ])
791 .await?
792 {
793 return Ok(());
794 }
795 // A repository renamed keeps its packages linked under its new
796 // name; one moved to another workspace leaves them behind, unlinked.
797 // Manage Actions access follows a rename, and lets go of a
798 // repository that left the package's workspace.
799 if g1t_kit::transfer::on_event(env, db, event, &[
800 "UPDATE packages SET repo_name = ?6 WHERE repo_id = ?5 AND workspace = ?3",
801 "UPDATE packages SET repo_id = NULL, repo_name = NULL WHERE repo_id = ?5 AND workspace <> ?3",
802 "UPDATE package_actions_access SET repo_name = ?6 WHERE repo_id = ?5",
803 "DELETE FROM package_actions_access WHERE repo_id = ?5 AND package_id IN (SELECT id FROM packages WHERE workspace <> ?3)",
804 ])
805 .await?
806 {
807 return Ok(());
808 }
809 if g1t_kit::lifecycle::on_purged(db, event, &[
810 "UPDATE packages SET repo_id = NULL, repo_name = NULL, visibility = 'private' WHERE repo_id = ?1",
811 "DELETE FROM package_actions_access WHERE repo_id = ?1",
812 ])
813 .await?
814 {
815 return Ok(());
816 }
817 // A team renamed or deleted: its roles on packages follow.
818 if matches!(event.kind.as_str(), "team.edited" | "team.deleted") {
819 let field = |snake: &str, camel: &str| {
820 event.data[snake].as_str().or_else(|| event.data[camel].as_str()).map(str::to_owned)
821 };
822 if let Some(team_id) = field("team_id", "teamId") {
823 let slug = field("team", "team");
824 let renamed = (event.kind == "team.edited").then_some(slug.as_deref()).flatten();
825 if event.kind == "team.deleted" || renamed.is_some() {
826 self.db.follow_team(&team_id, renamed).await?;
827 }
828 }
829 return Ok(());
830 }
831 if event.kind == "repo.visibility_changed" {
832 #[derive(serde::Deserialize)]
833 #[serde(rename_all = "camelCase")]
834 struct Changed {
835 repo_id: String,
836 is_private: bool,
837 }
838
839 let Ok(changed) = serde_json::from_value::<Changed>(event.data.clone()) else {
840 worker::console_error!("repo.visibility_changed {} could not be read", event.id);
841 return Ok(());
842 };
843 self.db.follow_visibility(&changed.repo_id, changed.is_private).await?;
844 for workspace in self.db.workspaces_linked_to(&changed.repo_id).await? {
845 self.db.measure(&workspace).await?;
846 }
847 }
848 Ok(())
849 }
850}
851
852#[event(fetch)]
853async fn fetch(mut request: Request, env: Env, ctx: Context) -> Result<Response> {
854 let packages = Packages::from_env(&env)?;
855 let Some(method) = rpc_method(&request) else {
856 if request.path().starts_with("/-/npm/") || request.path() == "/-/npm" {
857 return packages.npm(request, &ctx).await;
858 }
859 if request.path().starts_with("/-/composer/") {
860 return packages.composer(request, &ctx).await;
861 }
862 if request.path().starts_with("/-/cargo/") {
863 return packages.cargo(request, &ctx).await;
864 }
865 if request.path().starts_with("/-/maven/") {
866 return packages.maven(request, &ctx).await;
867 }
868 if request.path().starts_with("/-/nuget/") {
869 return packages.nuget(request, &ctx).await;
870 }
871 if request.path().starts_with("/-/rubygems/") {
872 return packages.rubygems(request, &ctx).await;
873 }
874 return packages.registry(request, &ctx).await;
875 };
876 let body: serde_json::Value = request.json().await?;
877 match method.as_str() {
878 "list_packages" => reply(&packages.list_packages(args(body)?).await?),
879 "get_package" => reply(&packages.get_package(args(body)?).await?),
880 "delete_version" => reply(&packages.delete_version(args(body)?).await?),
881 "delete_package" => reply(&packages.delete_package(args(body)?).await?),
882 "set_package" => reply(&packages.set_package(args(body)?).await?),
883 // A package's settings, deleted packages and versions: settings.rs.
884 "list_versions" => reply(&packages.list_versions(args(body)?).await?),
885 "get_version" => reply(&packages.get_version(args(body)?).await?),
886 "restore_version" => reply(&packages.restore_version(args(body)?).await?),
887 "restore_package" => reply(&packages.restore_package(args(body)?).await?),
888 "deleted_packages" => reply(&packages.deleted_packages(args(body)?).await?),
889 "package_settings" => reply(&packages.package_settings(args(body)?).await?),
890 "set_package_access" => reply(&packages.set_package_access(args(body)?).await?),
891 "remove_package_access" => reply(&packages.remove_package_access(args(body)?).await?),
892 "set_actions_access" => reply(&packages.set_actions_access(args(body)?).await?),
893 "remove_actions_access" => reply(&packages.remove_actions_access(args(body)?).await?),
894 // For billing: what a workspace's packages hold.
895 "storage_all" => reply(&packages.db.storage_all().await?),
896 // Read a repository's Composer package again now, as a push would.
897 "sync_composer" => {
898 let a: SyncComposerArgs = args(body)?;
899 reply(&packages.sync_composer(&a.repo_id).await?)
900 }
901 "storage" => {
902 let a: StorageArgs = args(body)?;
903 let (public_bytes, private_bytes) = packages.db.storage(&a.workspace.to_lowercase()).await?;
904 reply(&PackageStorage { public_bytes, private_bytes })
905 }
906 _ => Response::error("Unknown method", 404),
907 }
908}
909
910/// Every hour: expired uploads are let go, packages and versions deleted
911/// more than 30 days ago are purged, and blobs no version has used for a
912/// day are deleted from the store.
913#[event(scheduled)]
914async fn scheduled(_event: ScheduledEvent, env: Env, _ctx: ScheduleContext) {
915 let packages = match Packages::from_env(&env) {
916 Ok(packages) => packages,
917 Err(error) => {
918 worker::console_error!("packages: the sweep could not start: {error}");
919 return;
920 }
921 };
922 // Before the sweep, so the files of what is purged go with it a day on.
923 match packages.purge(now_ms()).await {
924 Ok((0, 0)) => {}
925 Ok((purged, versions)) => worker::console_log!("packages: purged {purged} deleted packages and {versions} deleted versions"),
926 Err(error) => worker::console_error!("packages: the purge failed: {error}"),
927 }
928 match packages.sweep().await {
929 Ok((0, 0)) => {}
930 Ok((uploads, blobs)) => worker::console_log!("packages: let go of {uploads} uploads and {blobs} blobs"),
931 Err(error) => worker::console_error!("packages: the sweep failed: {error}"),
932 }
933 // Repositories that had a composer.json before the registry did.
934 match packages.composer_backfill().await {
935 Ok(0) => {}
936 Ok(found) => worker::console_log!("packages: found {found} Composer packages"),
937 Err(error) => worker::console_error!("packages: the Composer backfill failed: {error}"),
938 }
939}
940
941#[event(queue)]
942async fn queue(batch: MessageBatch<Event>, env: Env, _ctx: Context) -> Result<()> {
943 let packages = Packages::from_env(&env)?;
944 for message in batch.messages()? {
945 packages.on_event(&env, message.body()).await?;
946 }
947 Ok(())
948}
949
950#[cfg(test)]
951mod tests {
952 use super::*;
953 use serde_json::json;
954
955 fn event(kind: &str, data: serde_json::Value) -> Event {
956 Event {
957 id: "evt_1".into(),
958 kind: kind.into(),
959 source: "identity".into(),
960 time: "2026-10-06T00:00:00.000Z".into(),
961 repo_id: None,
962 actor: None,
963 data,
964 }
965 }
966
967 #[test]
968 fn downloads_are_counted_by_version_and_written_in_batches() {
969 let mut counts = (HashMap::new(), 0);
970 let version = |v: &str| ("pkg_1".to_owned(), Some(v.to_owned()));
971 // The first is written at once (nothing was written before).
972 let first = tally(&mut counts, version("ver_1"), 100_000).unwrap();
973 assert_eq!(first, vec![(version("ver_1"), 1)]);
974 // Then they are kept until the flush is due, each version apart.
975 assert!(tally(&mut counts, version("ver_1"), 101_000).is_none());
976 assert!(tally(&mut counts, version("ver_1"), 102_000).is_none());
977 assert!(tally(&mut counts, version("ver_2"), 103_000).is_none());
978 let mut due = tally(&mut counts, ("pkg_1".to_owned(), None), 100_000 + DOWNLOADS_FLUSH_MS).unwrap();
979 due.sort();
980 assert_eq!(due, vec![(("pkg_1".to_owned(), None), 1), (version("ver_1"), 2), (version("ver_2"), 1)]);
981 assert!(counts.0.is_empty());
982 }
983
984 fn package(deleted_at: Option<&str>, workspace_deleted_at: Option<&str>) -> PackageRow {
985 PackageRow {
986 id: "pkg_1".into(),
987 workspace: "acme".into(),
988 ecosystem: "npm".into(),
989 name: "web".into(),
990 repo_id: None,
991 repo_name: None,
992 visibility: "private".into(),
993 description: None,
994 created_by: "usr_1".into(),
995 created_at: "2026-10-01T00:00:00.000Z".into(),
996 updated_at: "2026-10-01T00:00:00.000Z".into(),
997 downloads: 0,
998 workspace_deleted_at: workspace_deleted_at.map(str::to_owned),
999 inherit_access: 1,
1000 deleted_at: deleted_at.map(str::to_owned),
1001 deleted_by: deleted_at.map(|_| "ana".to_owned()),
1002 }
1003 }
1004
1005 #[test]
1006 fn a_deleted_package_is_hidden_and_keeps_its_name_until_the_purge() {
1007 let deleted = package(Some("2026-10-02T09:00:00.000Z"), None);
1008 assert!(deleted.hidden());
1009 let refusal = Packages::hidden_refusal(&deleted);
1010 assert!(refusal.contains("was deleted"), "{refusal}");
1011 assert!(refusal.contains("2026-11-01"), "kept 30 days: {refusal}");
1012 let workspace_gone = package(None, Some("2026-10-02T09:00:00.000Z"));
1013 assert!(workspace_gone.hidden());
1014 assert!(Packages::hidden_refusal(&workspace_gone).contains("workspace acme is deleted"));
1015 assert!(!package(None, None).hidden());
1016 }
1017
1018 #[test]
1019 fn the_purge_takes_what_was_deleted_more_than_thirty_days_ago() {
1020 let now = g1t_contracts::time::parse_rfc3339("2026-11-01T12:00:00.000Z").unwrap();
1021 let before = settings::purge_cutoff(now);
1022 assert_eq!(before, "2026-10-02T12:00:00.000Z");
1023 // Rows are purged when deleted_at < before: a day older goes, a
1024 // second newer stays restorable.
1025 assert!("2026-10-01T12:00:00.000Z" < before.as_str());
1026 assert!("2026-10-02T12:00:01.000Z" > before.as_str());
1027 assert!(settings::restorable("2026-10-02T12:00:01.000Z", now));
1028 assert!(!settings::restorable("2026-10-01T12:00:00.000Z", now));
1029 }
1030
1031 #[test]
1032 fn a_deleted_workspace_is_hidden_and_a_restored_one_shown() {
1033 let protected = g1t_contracts::identity::protected_names(Some("wsp_keep"));
1034 let now = "2026-10-06T12:00:00.000Z";
1035 let deleting = event("workspace.deleting", json!({ "workspaceId": "wsp_1", "slug": "Acme", "by": "ana", "purgeAfter": "x" }));
1036 assert_eq!(
1037 workspace_mark(&deleting, &protected, now),
1038 Some(WorkspaceMark::Hide { workspace: "acme".into(), at: now.into() })
1039 );
1040 let restored = event("workspace.restored", json!({ "workspaceId": "wsp_1", "slug": "acme" }));
1041 assert_eq!(workspace_mark(&restored, &protected, now), Some(WorkspaceMark::Show { workspace: "acme".into() }));
1042 }
1043
1044 #[test]
1045 fn protected_workspaces_and_other_events_are_left_alone() {
1046 let protected = g1t_contracts::identity::protected_names(Some("wsp_keep"));
1047 let now = "2026-10-06T12:00:00.000Z";
1048 let flagon = event("workspace.deleting", json!({ "workspaceId": "wsp_f", "slug": "flagon-io" }));
1049 assert_eq!(workspace_mark(&flagon, &protected, now), None, "always protected");
1050 let by_id = event("workspace.deleting", json!({ "workspaceId": "wsp_keep", "slug": "kept" }));
1051 assert_eq!(workspace_mark(&by_id, &protected, now), None, "protected by id");
1052 let purged = event("workspace.deleted", json!({ "workspaceId": "wsp_1", "slug": "acme" }));
1053 assert_eq!(workspace_mark(&purged, &protected, now), None, "the purge is handled apart");
1054 let nameless = event("workspace.deleting", json!({ "workspaceId": "wsp_1", "slug": "" }));
1055 assert_eq!(workspace_mark(&nameless, &protected, now), None);
1056 }
1057}