Skip to content
1,068 linesCodeBlameRaw
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, MessageExt, 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 let viewer: Viewer = 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 // An account that has not confirmed its email address signs in to
263 // nothing yet: its credentials are refused like wrong ones.
264 Ok(viewer.filter(|user| !user.awaits_confirmation()))
265 }
266
267 /// A repository of the workspace by name, whoever may see it.
268 async fn repo_by_name(&self, workspace: &str, name: &str) -> Result<Option<Repo>> {
269 if !g1t_contracts::is_valid_repo_name(name) {
270 return Ok(None);
271 }
272 let found: Outcome<Repo> = g1t_kit::call(
273 &self.repos,
274 "get",
275 &GetArgs {
276 path: RepoPath { namespace: workspace.to_owned(), name: name.to_owned() },
277 viewer: Some(User::system(workspace)),
278 },
279 )
280 .await?;
281 Ok(found.into_result().ok().filter(|repo| repo.fork_of.is_none()))
282 }
283
284 /// Who may do what to an image: the package's own settings, or, before
285 /// its first push, the repository it will be linked to.
286 async fn target(&self, name: &names::ImageName, found: Option<&PackageRow>) -> Result<TargetOf> {
287 if let Some(found) = found {
288 return Ok(TargetOf::package(found));
289 }
290 let repo = self.repo_by_name(&name.workspace, name.repo_name()).await?;
291 Ok(TargetOf::unmade(&name.workspace, &name.name, repo.map(|repo| (repo.id, repo.name, repo.is_private))))
292 }
293
294 /// What billing allows the workspace, kept for five minutes unless
295 /// `fresh`, and whether it is the kept answer. `None` when billing
296 /// cannot be asked: the push is then let through.
297 async fn allowance(&self, workspace: &str, fresh: bool) -> Option<(quota::Allowance, bool)> {
298 let now = now_ms();
299 let kept = ALLOWANCES.with(|kept| kept.borrow().get(workspace).copied());
300 if let Some((allowance, at)) = kept
301 && !fresh
302 && now.saturating_sub(at) < ALLOWANCE_TTL_MS
303 {
304 return Some((allowance, true));
305 }
306 let billing = self.env.service("BILLING").ok()?;
307 let asked: Result<quota::Allowance> = g1t_kit::call(
308 &billing,
309 "entitlements",
310 &g1t_contracts::billing::EntitlementsArgs { workspace: workspace.to_owned() },
311 )
312 .await;
313 match asked {
314 Ok(allowance) => {
315 ALLOWANCES.with(|kept| kept.borrow_mut().insert(workspace.to_owned(), (allowance, now)));
316 Some((allowance, false))
317 }
318 Err(error) => {
319 worker::console_error!("packages: billing could not be asked about {workspace}, letting the push through: {error}");
320 None
321 }
322 }
323 }
324
325 /// Why a push of these files (digest and size) into the package may
326 /// not be kept, if it may not: a free workspace past its free storage.
327 /// Files the workspace already holds add nothing.
328 pub(crate) async fn storage_refusal(&self, package: &PackageRow, files: &[(String, u64)]) -> Result<Option<String>> {
329 if !self.storage_limits || files.is_empty() {
330 return Ok(None);
331 }
332 let digests: Vec<String> = files.iter().map(|(digest, _)| digest.clone()).collect();
333 let held = self.db.held(&package.workspace, &digests).await?;
334 let mut seen = std::collections::HashSet::new();
335 let adding: u64 = files
336 .iter()
337 .filter(|(digest, _)| !held.contains(digest) && seen.insert(digest.clone()))
338 .map(|(_, size)| size)
339 .sum();
340 if adding == 0 {
341 return Ok(None);
342 }
343 let Some((allowance, kept)) = self.allowance(&package.workspace, false).await else {
344 return Ok(None);
345 };
346 let (public_bytes, private_bytes) = self.db.storage(&package.workspace).await?;
347 let public = package.public();
348 let used = if public { public_bytes } else { private_bytes };
349 let mut refused = quota::decide(&allowance, public, used, adding).err();
350 // A refusal from the kept answer is checked with billing again: the
351 // workspace may have just added a plan, and must not wait minutes
352 // for the push to go through.
353 if refused.is_some() && kept {
354 refused = match self.allowance(&package.workspace, true).await {
355 Some((allowance, _)) => quota::decide(&allowance, public, used, adding).err(),
356 None => None,
357 };
358 }
359 Ok(refused.map(|refusal| quota::message(&package.workspace, &refusal)))
360 }
361
362 fn count_download(&self, package_id: &str, ctx: &Context) {
363 self.count_downloads(package_id, None, ctx);
364 }
365
366 /// A download of one version, counted for it and its package.
367 fn count_version_download(&self, package_id: &str, version_id: &str, ctx: &Context) {
368 self.count_downloads(package_id, Some(version_id), ctx);
369 }
370
371 fn count_downloads(&self, package_id: &str, version_id: Option<&str>, ctx: &Context) {
372 let key = (package_id.to_owned(), version_id.map(str::to_owned));
373 let due = DOWNLOADS.with(|counts| tally(&mut counts.borrow_mut(), key, now_ms()));
374 if let Some(due) = due {
375 let env = self.env.clone();
376 ctx.wait_until(async move {
377 let written = match env.d1("DB") {
378 Ok(db) => Db { db }.add_downloads(&due).await,
379 Err(error) => Err(error),
380 };
381 if let Err(error) = written {
382 worker::console_error!("packages: downloads not counted: {error}");
383 }
384 });
385 }
386 }
387
388 async fn announce(&self, kind: &'static str, package: &PackageRow, data: PackageEvent, caller: &Caller) {
389 let event = NewEvent {
390 kind,
391 source: SOURCE,
392 repo_id: package.repo_id.clone(),
393 actor: caller.actor.as_ref().map(|actor| actor.actor_id.clone()),
394 data,
395 };
396 let published: Result<serde_json::Value> = g1t_kit::call(&self.events, "publish", &Publish { events: vec![event] }).await;
397 if let Err(error) = published {
398 worker::console_error!("packages: {kind} not published: {error}");
399 }
400 }
401
402 async fn audit(&self, caller: &Caller, action: &str, package: &PackageRow, path: Option<&str>, surface: Option<Surface>) {
403 self.audit_with(caller, action, package, path, surface, None).await;
404 }
405
406 /// An audit entry, with what changed in words (`message`), as identity
407 /// writes its access entries.
408 async fn audit_with(
409 &self,
410 caller: &Caller,
411 action: &str,
412 package: &PackageRow,
413 path: Option<&str>,
414 surface: Option<Surface>,
415 message: Option<String>,
416 ) {
417 let Some(actor) = caller.actor.clone() else {
418 return;
419 };
420 let target = AuditTarget {
421 workspace: package.workspace.clone(),
422 repo: package.repo_name.as_ref().map(|name| format!("{}/{name}", package.workspace)),
423 path: Some(path.map_or_else(|| format!("{}:{}/{}", package.ecosystem, package.workspace, package.name), str::to_owned)),
424 ..AuditTarget::default()
425 };
426 let mut entry = NewAuditEntry::new(
427 actor,
428 action,
429 surface.unwrap_or(Surface::Registry),
430 target,
431 &Decision::allow("packages"),
432 new_id("req", now_ms()),
433 );
434 entry.result = Some("ok".to_owned());
435 entry.message = message;
436 let recorded: Result<u32> = g1t_kit::call(&self.events, "audit_record", &RecordAuditArgs { entries: vec![entry] }).await;
437 if let Err(error) = recorded {
438 worker::console_error!("packages: audit entry not recorded: {error}");
439 }
440 }
441
442 fn summary(&self, row: &db::ListedRow) -> PackageSummary {
443 let p = &row.package;
444 PackageSummary {
445 id: p.id.clone(),
446 workspace: p.workspace.clone(),
447 ecosystem: Ecosystem::parse(&p.ecosystem).unwrap_or(Ecosystem::Container),
448 name: p.name.clone(),
449 address: match p.ecosystem.as_str() {
450 "npm" => format!("{}/-/npm/@{}/{}", self.host, p.workspace, p.name),
451 "composer" => format!("{}/-/composer/{}/{}", self.host, p.workspace, p.name),
452 "cargo" => format!("{}/-/cargo/{}/{}", self.host, p.workspace, p.name),
453 "maven" => format!("{}/-/maven/{}/{}", self.host, p.workspace, p.name),
454 "nuget" => format!("{}/-/nuget/{}/{}", self.host, p.workspace, p.name),
455 "rubygems" => format!("{}/-/rubygems/{}/{}", self.host, p.workspace, p.name),
456 _ => format!("{}/{}/{}", self.host, p.workspace, p.name),
457 },
458 visibility: Visibility::parse(&p.visibility),
459 repo: p.repo_id.as_ref().map(|id| LinkedRepo {
460 id: id.clone(),
461 namespace: p.workspace.clone(),
462 name: p.repo_name.clone().unwrap_or_default(),
463 }),
464 description: p.description.clone(),
465 versions: row.version_count,
466 latest: db::latest_shown(row),
467 size: row.bytes,
468 downloads: p.downloads,
469 created_at: p.created_at.clone(),
470 updated_at: p.updated_at.clone(),
471 inherit_access: p.inherits(),
472 deleted_at: p.deleted_at.clone(),
473 deleted_by: p.deleted_by.clone(),
474 purge_at: p.deleted_at.as_deref().and_then(settings::purge_at),
475 }
476 }
477
478 async fn listed(&self, workspace: &str, ecosystem: Ecosystem, name: &str) -> Result<Option<db::ListedRow>> {
479 let rows = self.db.list(workspace, Some(ecosystem.as_str()), None, None, LIST_LIMIT).await?;
480 if let Some(row) = rows.into_iter().find(|row| row.package.name == name) {
481 return Ok(Some(row));
482 }
483 // Past the first page: read it alone.
484 Ok(self.db.package(workspace, ecosystem.as_str(), name).await?.filter(|package| !package.hidden()).map(|package| db::ListedRow {
485 package,
486 version_count: 0,
487 bytes: 0,
488 latest_tag: None,
489 latest_tag_version: None,
490 latest_version: None,
491 }))
492 }
493
494 async fn list_packages(&self, a: ListPackagesArgs) -> Result<Outcome<Vec<PackageSummary>>> {
495 let workspace = a.workspace.to_lowercase();
496 let rows = self
497 .db
498 .list(&workspace, a.ecosystem.map(Ecosystem::as_str), a.repo_id.as_deref(), a.query.as_deref(), LIST_LIMIT)
499 .await?;
500 let packages: Vec<&PackageRow> = rows.iter().map(|row| &row.package).collect();
501 let may = self.may_all(a.viewer.as_ref(), &packages, Action::Pull).await?;
502 let visible = rows
503 .iter()
504 .zip(may)
505 .filter(|(_, may)| *may)
506 .map(|(row, _)| self.summary(row))
507 .collect();
508 Ok(Outcome::Ok(visible))
509 }
510
511 async fn get_package(&self, a: GetPackageArgs) -> Result<Outcome<PackageDetail>> {
512 let workspace = a.workspace.to_lowercase();
513 let Some(row) = self.listed(&workspace, a.ecosystem, &a.name).await? else {
514 return Ok(not_found());
515 };
516 let mut target = TargetOf::package(&row.package);
517 let permissions = self.permissions(a.viewer.as_ref(), &mut target).await?;
518 if !permissions.pull {
519 return Ok(not_found());
520 }
521 let tags = self.db.tags(&row.package.id).await?;
522 let versions = self
523 .db
524 .versions(&row.package.id, VERSIONS_SHOWN)
525 .await?
526 .into_iter()
527 .map(|version| settings::version_of(&row.package, version, &tags))
528 .collect();
529 // The README its page shows: npm's, from the latest version.
530 let readme = match self.db.readme_digest(&row.package.id).await?.and_then(|d| digest::Digest::parse(&d)) {
531 Some(digest) => match self.db.blob(&digest).await? {
532 Some(blob) => self.store.read(&blob.object_key).await?.map(|b| String::from_utf8_lossy(&b).into_owned()),
533 None => None,
534 },
535 None => None,
536 };
537 Ok(Outcome::Ok(PackageDetail {
538 readme,
539 package: self.summary(&row),
540 versions,
541 tags: tags
542 .into_iter()
543 .map(|tag| PackageTag { tag: tag.tag, digest: tag.digest, updated_at: tag.updated_at })
544 .collect(),
545 permissions,
546 }))
547 }
548
549 /// The package `actor` asks to change, if they may `action` it.
550 pub(crate) async fn for_change(&self, actor: &User, workspace: &str, ecosystem: Ecosystem, name: &str, action: Action) -> Result<Outcome<PackageRow>> {
551 let Some(package) = self.db.package(&workspace.to_lowercase(), ecosystem.as_str(), name).await?.filter(|p| !p.hidden()) else {
552 return Ok(not_found());
553 };
554 match self.allowed(actor, &package, action).await? {
555 Outcome::Ok(()) => Ok(Outcome::Ok(package)),
556 Outcome::Fail(failure) => Ok(Outcome::Fail(failure)),
557 }
558 }
559
560 /// Whether `actor` may `action` the package: not found for one who may
561 /// not pull it, forbidden with the reason for one who may.
562 pub(crate) async fn allowed(&self, actor: &User, package: &PackageRow, action: Action) -> Result<Outcome<()>> {
563 let mut target = TargetOf::package(package);
564 let decision = self.decide(Some(actor), &mut target, action).await?;
565 if decision.allowed {
566 return Ok(Outcome::Ok(()));
567 }
568 if !self.decide(Some(actor), &mut target, Action::Pull).await?.allowed {
569 return Ok(not_found());
570 }
571 Ok(Outcome::fail(FailureCode::Forbidden, decision.reason.unwrap_or_else(|| "Not allowed.".to_owned())))
572 }
573
574 async fn delete_version(&self, a: DeleteVersionArgs) -> Result<Outcome<()>> {
575 let package = match self.for_change(&a.actor, &a.workspace, a.ecosystem, &a.name, Action::Delete).await? {
576 Outcome::Ok(package) => package,
577 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
578 };
579 if package.ecosystem == "composer" {
580 return Ok(Outcome::fail(
581 FailureCode::Invalid,
582 "A Composer package's versions are its repository's tags and branches: delete the tag or branch instead.",
583 ));
584 }
585 let Some(version) = self.db.find_version(&package.id, &a.version).await? else {
586 return Ok(Outcome::fail(FailureCode::NotFound, "Version not found."));
587 };
588 let caller = Caller::of(Some(&a.actor));
589 self.remove_version_from(&package, &version, &caller, a.surface).await?;
590 Ok(Outcome::Ok(()))
591 }
592
593 async fn delete_package(&self, a: DeletePackageArgs) -> Result<Outcome<()>> {
594 let package = match self.for_change(&a.actor, &a.workspace, a.ecosystem, &a.name, Action::Delete).await? {
595 Outcome::Ok(package) => package,
596 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
597 };
598 let caller = Caller::of(Some(&a.actor));
599 self.remove_package(&package, &caller, a.surface).await?;
600 Ok(Outcome::Ok(()))
601 }
602
603 /// Deletes a package: hidden at once, restorable for 30 days, its name
604 /// kept until then. Its files go with the purge.
605 pub(crate) async fn remove_package(&self, package: &PackageRow, caller: &Caller, surface: Option<Surface>) -> Result<()> {
606 let by = caller.actor.as_ref().map_or("", |actor| actor.actor.as_str());
607 self.db.soft_delete_package(&package.id, by, now_ms()).await?;
608 self.db.measure(&package.workspace).await?;
609 self.announce("package.deleted", package, self.event_of(package), caller).await;
610 self.audit_with(
611 caller,
612 "package.delete",
613 package,
614 None,
615 surface,
616 Some(format!("Deleted the package; it can be restored for {RESTORE_DAYS} days")),
617 )
618 .await;
619 Ok(())
620 }
621
622 async fn set_package(&self, a: SetPackageArgs) -> Result<Outcome<PackageSummary>> {
623 let package = match self.for_change(&a.actor, &a.workspace, a.ecosystem, &a.name, Action::Admin).await? {
624 Outcome::Ok(package) => package,
625 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
626 };
627 let caller = Caller::of(Some(&a.actor));
628 let now = now_ms();
629 let before = package.visibility.clone();
630 let linked = if a.unlink { false } else { a.link.is_some() || package.repo_id.is_some() };
631 if a.visibility.is_some_and(|visibility| visibility.as_str() != package.visibility) && linked {
632 return Ok(Outcome::fail(
633 FailureCode::Invalid,
634 "A linked package has its repository's visibility. Change the repository's, or unlink the package.",
635 ));
636 }
637 if a.inherit_access.is_some() && !linked {
638 return Ok(Outcome::fail(
639 FailureCode::Invalid,
640 "Only a package linked to a repository inherits access. Link it to a repository first.",
641 ));
642 }
643 let path = |repo: &str| format!("{}/{repo}", package.workspace);
644 if a.unlink {
645 if let Some(name) = package.repo_name.as_deref() {
646 self.db.set_link(&package.id, None, "private", now).await?;
647 self.audit_with(&caller, "package.unlinked", &package, None, a.surface, Some(format!("Unlinked from {}", path(name)))).await;
648 }
649 } else if let Some(link) = a.link.as_deref() {
650 let link = link.trim().trim_end_matches(".git");
651 let name = match link.split_once('/') {
652 Some((owner, name)) if owner.eq_ignore_ascii_case(&package.workspace) => name,
653 Some(_) => {
654 return Ok(Outcome::fail(
655 FailureCode::Invalid,
656 format!("A package can only be linked to a repository of its own workspace, {}.", package.workspace),
657 ));
658 }
659 None => link,
660 };
661 let Some(repo) = self.repo_by_name(&package.workspace, name).await? else {
662 return Ok(Outcome::fail(FailureCode::NotFound, format!("There is no repository {}.", path(name))));
663 };
664 // Linking hands the package to the repository's roles: only
665 // someone who administers that repository may.
666 let role = g1t_contracts::access::permission(Some(&a.actor), &repo);
667 if role < Some(g1t_contracts::access::RepoRole::Admin) {
668 return Ok(Outcome::fail(
669 FailureCode::Forbidden,
670 format!("You need the Admin role on {}/{} to link a package to it.", repo.namespace, repo.name),
671 ));
672 }
673 if package.repo_id.as_deref() != Some(repo.id.as_str()) {
674 let visibility = if repo.is_private { "private" } else { "public" };
675 self.db.set_link(&package.id, Some((&repo.id, &repo.name)), visibility, now).await?;
676 self.audit_with(&caller, "package.linked", &package, None, a.surface, Some(format!("Linked to {}", path(&repo.name)))).await;
677 }
678 }
679 if let Some(visibility) = a.visibility
680 && !linked
681 && visibility.as_str() != package.visibility
682 {
683 self.db.set_visibility(&package.id, visibility.as_str(), now).await?;
684 }
685 if let Some(inherit) = a.inherit_access
686 && inherit != package.inherits()
687 {
688 self.db.set_inherit(&package.id, inherit, now).await?;
689 let message = if inherit {
690 "Turned on inheriting access from the linked repository"
691 } else {
692 "Turned off inheriting access from the linked repository"
693 };
694 self.audit_with(&caller, "package.inherit_access_changed", &package, None, a.surface, Some(message.to_owned())).await;
695 }
696 let Some(row) = self.listed(&package.workspace, a.ecosystem, &package.name).await? else {
697 return Ok(not_found());
698 };
699 if row.package.visibility != before {
700 self.db.measure(&package.workspace).await?;
701 let event = PackageEvent { visibility: Some(row.package.visibility.clone()), ..self.event_of(&row.package) };
702 self.announce("package.visibility_changed", &row.package, event, &caller).await;
703 self.audit_with(
704 &caller,
705 "package.visibility_changed",
706 &row.package,
707 None,
708 a.surface,
709 Some(format!("Made the package {}", row.package.visibility)),
710 )
711 .await;
712 }
713 Ok(Outcome::Ok(self.summary(&row)))
714 }
715
716 /// Lets go of expired uploads, and of blobs no version has used for a day.
717 async fn sweep(&self) -> Result<(u32, u32)> {
718 let now = now_ms();
719 let mut uploads = 0;
720 for row in self.db.expired_uploads(now, SWEEP_BATCH).await? {
721 if let Some(progress) = row.progress()
722 && let Err(error) = upload::abort(&self.store, &progress).await
723 {
724 worker::console_error!("packages: upload {} not aborted: {error}", row.id);
725 }
726 self.db.delete_upload(&row.id).await?;
727 uploads += 1;
728 }
729 let mut blobs = 0;
730 for blob in self.db.unused_blobs(now, SWEEP_BATCH).await? {
731 self.store.delete(&blob.object_key).await?;
732 self.db.forget_blob(&blob.digest).await?;
733 blobs += 1;
734 }
735 Ok((uploads, blobs))
736 }
737
738 /// Follows what happens elsewhere: workspaces renamed and deleted,
739 /// repositories that change visibility, are renamed, move or go.
740 async fn on_event(&self, env: &Env, event: &Event) -> Result<()> {
741 // Composer packages are read from repositories: again when one is
742 // pushed to (its default branch may have gained a composer.json),
743 // restored, renamed or moved; gone when it is deleted. A failure is
744 // logged, not retried with the batch: the next push reads it again.
745 let repo_id = event.repo_id.clone().or_else(|| event.data["repoId"].as_str().map(str::to_owned));
746 if let Some(repo_id) = repo_id.as_deref() {
747 let synced = match event.kind.as_str() {
748 "git.push" => {
749 let known = self.db.package_for_repo(repo_id, "composer").await?.is_some();
750 if known || event.data["defaultBranch"].as_bool() == Some(true) {
751 Some(self.sync_composer(repo_id).await.map(|_| ()))
752 } else {
753 None
754 }
755 }
756 "repo.deleted" => Some(self.composer_repo_gone(repo_id).await),
757 "repo.restored" | "repo.renamed" | "repo.transferred" => Some(self.sync_composer(repo_id).await.map(|_| ())),
758 _ => None,
759 };
760 if let Some(Err(error)) = synced {
761 worker::console_error!("packages: composer {} for {repo_id}: {error}", event.kind);
762 }
763 if event.kind == "git.push" || event.kind == "repo.deleted" || event.kind == "repo.restored" {
764 return Ok(());
765 }
766 }
767 let db = &self.db.db;
768 let protected = g1t_contracts::identity::protected_names(Some(&store::var(env, "PROTECTED_WORKSPACES")));
769 match workspace_mark(event, &protected, &g1t_contracts::time::rfc3339(now_ms())) {
770 Some(WorkspaceMark::Hide { workspace, at }) => return self.db.mark_workspace(&workspace, Some(&at)).await,
771 Some(WorkspaceMark::Show { workspace }) => return self.db.mark_workspace(&workspace, None).await,
772 None => {}
773 }
774 if g1t_kit::rename::on_event(env, db, event, &[
775 "UPDATE packages SET workspace = ?1 WHERE workspace = ?2",
776 "UPDATE uploads SET workspace = ?1 WHERE workspace = ?2",
777 "UPDATE workspace_blobs SET workspace = ?1 WHERE workspace = ?2",
778 ])
779 .await?
780 {
781 return Ok(());
782 }
783 if g1t_kit::deleted::on_event(db, event, &[
784 "DELETE FROM tags WHERE package_id IN (SELECT id FROM packages WHERE workspace = ?1)",
785 "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)",
786 "DELETE FROM versions WHERE package_id IN (SELECT id FROM packages WHERE workspace = ?1)",
787 "DELETE FROM package_blobs WHERE package_id IN (SELECT id FROM packages WHERE workspace = ?1)",
788 "DELETE FROM package_access WHERE package_id IN (SELECT id FROM packages WHERE workspace = ?1)",
789 "DELETE FROM package_actions_access WHERE package_id IN (SELECT id FROM packages WHERE workspace = ?1)",
790 "DELETE FROM uploads WHERE workspace = ?1",
791 "DELETE FROM workspace_blobs WHERE workspace = ?1",
792 "DELETE FROM packages WHERE workspace = ?1",
793 ])
794 .await?
795 {
796 return Ok(());
797 }
798 // A repository renamed keeps its packages linked under its new
799 // name; one moved to another workspace leaves them behind, unlinked.
800 // Manage Actions access follows a rename, and lets go of a
801 // repository that left the package's workspace.
802 if g1t_kit::transfer::on_event(env, db, event, &[
803 "UPDATE packages SET repo_name = ?6 WHERE repo_id = ?5 AND workspace = ?3",
804 "UPDATE packages SET repo_id = NULL, repo_name = NULL WHERE repo_id = ?5 AND workspace <> ?3",
805 "UPDATE package_actions_access SET repo_name = ?6 WHERE repo_id = ?5",
806 "DELETE FROM package_actions_access WHERE repo_id = ?5 AND package_id IN (SELECT id FROM packages WHERE workspace <> ?3)",
807 ])
808 .await?
809 {
810 return Ok(());
811 }
812 if g1t_kit::lifecycle::on_purged(db, event, &[
813 "UPDATE packages SET repo_id = NULL, repo_name = NULL, visibility = 'private' WHERE repo_id = ?1",
814 "DELETE FROM package_actions_access WHERE repo_id = ?1",
815 ])
816 .await?
817 {
818 return Ok(());
819 }
820 // A team renamed or deleted: its roles on packages follow.
821 if matches!(event.kind.as_str(), "team.edited" | "team.deleted") {
822 let field = |snake: &str, camel: &str| {
823 event.data[snake].as_str().or_else(|| event.data[camel].as_str()).map(str::to_owned)
824 };
825 if let Some(team_id) = field("team_id", "teamId") {
826 let slug = field("team", "team");
827 let renamed = (event.kind == "team.edited").then_some(slug.as_deref()).flatten();
828 if event.kind == "team.deleted" || renamed.is_some() {
829 self.db.follow_team(&team_id, renamed).await?;
830 }
831 }
832 return Ok(());
833 }
834 if event.kind == "repo.visibility_changed" {
835 #[derive(serde::Deserialize)]
836 #[serde(rename_all = "camelCase")]
837 struct Changed {
838 repo_id: String,
839 is_private: bool,
840 }
841
842 let Ok(changed) = serde_json::from_value::<Changed>(event.data.clone()) else {
843 worker::console_error!("repo.visibility_changed {} could not be read", event.id);
844 return Ok(());
845 };
846 self.db.follow_visibility(&changed.repo_id, changed.is_private).await?;
847 for workspace in self.db.workspaces_linked_to(&changed.repo_id).await? {
848 self.db.measure(&workspace).await?;
849 }
850 }
851 Ok(())
852 }
853}
854
855#[event(fetch)]
856async fn fetch(mut request: Request, env: Env, ctx: Context) -> Result<Response> {
857 let packages = Packages::from_env(&env)?;
858 let Some(method) = rpc_method(&request) else {
859 if request.path().starts_with("/-/npm/") || request.path() == "/-/npm" {
860 return packages.npm(request, &ctx).await;
861 }
862 if request.path().starts_with("/-/composer/") {
863 return packages.composer(request, &ctx).await;
864 }
865 if request.path().starts_with("/-/cargo/") {
866 return packages.cargo(request, &ctx).await;
867 }
868 if request.path().starts_with("/-/maven/") {
869 return packages.maven(request, &ctx).await;
870 }
871 if request.path().starts_with("/-/nuget/") {
872 return packages.nuget(request, &ctx).await;
873 }
874 if request.path().starts_with("/-/rubygems/") {
875 return packages.rubygems(request, &ctx).await;
876 }
877 return packages.registry(request, &ctx).await;
878 };
879 let body: serde_json::Value = request.json().await?;
880 match method.as_str() {
881 "list_packages" => reply(&packages.list_packages(args(body)?).await?),
882 "get_package" => reply(&packages.get_package(args(body)?).await?),
883 "delete_version" => reply(&packages.delete_version(args(body)?).await?),
884 "delete_package" => reply(&packages.delete_package(args(body)?).await?),
885 "set_package" => reply(&packages.set_package(args(body)?).await?),
886 // A package's settings, deleted packages and versions: settings.rs.
887 "list_versions" => reply(&packages.list_versions(args(body)?).await?),
888 "get_version" => reply(&packages.get_version(args(body)?).await?),
889 "restore_version" => reply(&packages.restore_version(args(body)?).await?),
890 "restore_package" => reply(&packages.restore_package(args(body)?).await?),
891 "deleted_packages" => reply(&packages.deleted_packages(args(body)?).await?),
892 "package_settings" => reply(&packages.package_settings(args(body)?).await?),
893 "set_package_access" => reply(&packages.set_package_access(args(body)?).await?),
894 "remove_package_access" => reply(&packages.remove_package_access(args(body)?).await?),
895 "set_actions_access" => reply(&packages.set_actions_access(args(body)?).await?),
896 "remove_actions_access" => reply(&packages.remove_actions_access(args(body)?).await?),
897 // For billing: what a workspace's packages hold.
898 "storage_all" => reply(&packages.db.storage_all().await?),
899 // Read a repository's Composer package again now, as a push would.
900 "sync_composer" => {
901 let a: SyncComposerArgs = args(body)?;
902 reply(&packages.sync_composer(&a.repo_id).await?)
903 }
904 "storage" => {
905 let a: StorageArgs = args(body)?;
906 let (public_bytes, private_bytes) = packages.db.storage(&a.workspace.to_lowercase()).await?;
907 reply(&PackageStorage { public_bytes, private_bytes })
908 }
909 _ => Response::error("Unknown method", 404),
910 }
911}
912
913/// Every hour: expired uploads are let go, packages and versions deleted
914/// more than 30 days ago are purged, and blobs no version has used for a
915/// day are deleted from the store.
916#[event(scheduled)]
917async fn scheduled(_event: ScheduledEvent, env: Env, _ctx: ScheduleContext) {
918 let packages = match Packages::from_env(&env) {
919 Ok(packages) => packages,
920 Err(error) => {
921 worker::console_error!("packages: the sweep could not start: {error}");
922 return;
923 }
924 };
925 // Before the sweep, so the files of what is purged go with it a day on.
926 match packages.purge(now_ms()).await {
927 Ok((0, 0)) => {}
928 Ok((purged, versions)) => worker::console_log!("packages: purged {purged} deleted packages and {versions} deleted versions"),
929 Err(error) => worker::console_error!("packages: the purge failed: {error}"),
930 }
931 match packages.sweep().await {
932 Ok((0, 0)) => {}
933 Ok((uploads, blobs)) => worker::console_log!("packages: let go of {uploads} uploads and {blobs} blobs"),
934 Err(error) => worker::console_error!("packages: the sweep failed: {error}"),
935 }
936 // Repositories that had a composer.json before the registry did.
937 match packages.composer_backfill().await {
938 Ok(0) => {}
939 Ok(found) => worker::console_log!("packages: found {found} Composer packages"),
940 Err(error) => worker::console_error!("packages: the Composer backfill failed: {error}"),
941 }
942}
943
944#[event(queue)]
945async fn queue(batch: MessageBatch<Event>, env: Env, _ctx: Context) -> Result<()> {
946 let packages = Packages::from_env(&env)?;
947 // Each event is acknowledged or retried on its own, so one that fails
948 // does not run the rest of its batch again.
949 for message in batch.messages()? {
950 match packages.on_event(&env, message.body()).await {
951 Ok(_) => message.ack(),
952 Err(error) => {
953 worker::console_error!("packages: event {} failed: {error}", message.body().id);
954 message.retry();
955 }
956 }
957 }
958 Ok(())
959}
960
961#[cfg(test)]
962mod tests {
963 use super::*;
964 use serde_json::json;
965
966 fn event(kind: &str, data: serde_json::Value) -> Event {
967 Event {
968 id: "evt_1".into(),
969 kind: kind.into(),
970 source: "identity".into(),
971 time: "2026-10-06T00:00:00.000Z".into(),
972 repo_id: None,
973 actor: None,
974 data,
975 }
976 }
977
978 #[test]
979 fn downloads_are_counted_by_version_and_written_in_batches() {
980 let mut counts = (HashMap::new(), 0);
981 let version = |v: &str| ("pkg_1".to_owned(), Some(v.to_owned()));
982 // The first is written at once (nothing was written before).
983 let first = tally(&mut counts, version("ver_1"), 100_000).unwrap();
984 assert_eq!(first, vec![(version("ver_1"), 1)]);
985 // Then they are kept until the flush is due, each version apart.
986 assert!(tally(&mut counts, version("ver_1"), 101_000).is_none());
987 assert!(tally(&mut counts, version("ver_1"), 102_000).is_none());
988 assert!(tally(&mut counts, version("ver_2"), 103_000).is_none());
989 let mut due = tally(&mut counts, ("pkg_1".to_owned(), None), 100_000 + DOWNLOADS_FLUSH_MS).unwrap();
990 due.sort();
991 assert_eq!(due, vec![(("pkg_1".to_owned(), None), 1), (version("ver_1"), 2), (version("ver_2"), 1)]);
992 assert!(counts.0.is_empty());
993 }
994
995 fn package(deleted_at: Option<&str>, workspace_deleted_at: Option<&str>) -> PackageRow {
996 PackageRow {
997 id: "pkg_1".into(),
998 workspace: "acme".into(),
999 ecosystem: "npm".into(),
1000 name: "web".into(),
1001 repo_id: None,
1002 repo_name: None,
1003 visibility: "private".into(),
1004 description: None,
1005 created_by: "usr_1".into(),
1006 created_at: "2026-10-01T00:00:00.000Z".into(),
1007 updated_at: "2026-10-01T00:00:00.000Z".into(),
1008 downloads: 0,
1009 workspace_deleted_at: workspace_deleted_at.map(str::to_owned),
1010 inherit_access: 1,
1011 deleted_at: deleted_at.map(str::to_owned),
1012 deleted_by: deleted_at.map(|_| "ana".to_owned()),
1013 }
1014 }
1015
1016 #[test]
1017 fn a_deleted_package_is_hidden_and_keeps_its_name_until_the_purge() {
1018 let deleted = package(Some("2026-10-02T09:00:00.000Z"), None);
1019 assert!(deleted.hidden());
1020 let refusal = Packages::hidden_refusal(&deleted);
1021 assert!(refusal.contains("was deleted"), "{refusal}");
1022 assert!(refusal.contains("2026-11-01"), "kept 30 days: {refusal}");
1023 let workspace_gone = package(None, Some("2026-10-02T09:00:00.000Z"));
1024 assert!(workspace_gone.hidden());
1025 assert!(Packages::hidden_refusal(&workspace_gone).contains("workspace acme is deleted"));
1026 assert!(!package(None, None).hidden());
1027 }
1028
1029 #[test]
1030 fn the_purge_takes_what_was_deleted_more_than_thirty_days_ago() {
1031 let now = g1t_contracts::time::parse_rfc3339("2026-11-01T12:00:00.000Z").unwrap();
1032 let before = settings::purge_cutoff(now);
1033 assert_eq!(before, "2026-10-02T12:00:00.000Z");
1034 // Rows are purged when deleted_at < before: a day older goes, a
1035 // second newer stays restorable.
1036 assert!("2026-10-01T12:00:00.000Z" < before.as_str());
1037 assert!("2026-10-02T12:00:01.000Z" > before.as_str());
1038 assert!(settings::restorable("2026-10-02T12:00:01.000Z", now));
1039 assert!(!settings::restorable("2026-10-01T12:00:00.000Z", now));
1040 }
1041
1042 #[test]
1043 fn a_deleted_workspace_is_hidden_and_a_restored_one_shown() {
1044 let protected = g1t_contracts::identity::protected_names(Some("wsp_keep"));
1045 let now = "2026-10-06T12:00:00.000Z";
1046 let deleting = event("workspace.deleting", json!({ "workspaceId": "wsp_1", "slug": "Acme", "by": "ana", "purgeAfter": "x" }));
1047 assert_eq!(
1048 workspace_mark(&deleting, &protected, now),
1049 Some(WorkspaceMark::Hide { workspace: "acme".into(), at: now.into() })
1050 );
1051 let restored = event("workspace.restored", json!({ "workspaceId": "wsp_1", "slug": "acme" }));
1052 assert_eq!(workspace_mark(&restored, &protected, now), Some(WorkspaceMark::Show { workspace: "acme".into() }));
1053 }
1054
1055 #[test]
1056 fn protected_workspaces_and_other_events_are_left_alone() {
1057 let protected = g1t_contracts::identity::protected_names(Some("wsp_keep"));
1058 let now = "2026-10-06T12:00:00.000Z";
1059 let flagon = event("workspace.deleting", json!({ "workspaceId": "wsp_f", "slug": "flagon-io" }));
1060 assert_eq!(workspace_mark(&flagon, &protected, now), None, "always protected");
1061 let by_id = event("workspace.deleting", json!({ "workspaceId": "wsp_keep", "slug": "kept" }));
1062 assert_eq!(workspace_mark(&by_id, &protected, now), None, "protected by id");
1063 let purged = event("workspace.deleted", json!({ "workspaceId": "wsp_1", "slug": "acme" }));
1064 assert_eq!(workspace_mark(&purged, &protected, now), None, "the purge is handled apart");
1065 let nameless = event("workspace.deleting", json!({ "workspaceId": "wsp_1", "slug": "" }));
1066 assert_eq!(workspace_mark(&nameless, &protected, now), None);
1067 }
1068}