g1t/services/packages/src/lib.rs

835 lines36,687 bytesCodeBlame
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
56pub(crate) const SOURCE: &str = "packages";
57/// How large a request body may be: Cloudflare's limit on the zone's plan.
58const DEFAULT_MAX_REQUEST_BYTES: u64 = 100_000_000;
59/// How many packages a listing shows.
60const LIST_LIMIT: u32 = 200;
61const VERSIONS_SHOWN: u32 = 200;
62/// How much one sweep lets go of.
63const SWEEP_BATCH: u32 = 200;
64
65thread_local! {
66 /// Pulls counted since the last write, by package (and by version,
67 /// where it is counted too): written at most every few seconds, so a
68 /// busy image costs one write, not one a pull. What an isolate holds
69 /// when it goes away is lost: the count is approximate.
70 static DOWNLOADS: RefCell<(HashMap<DownloadKey, u64>, u64)> = RefCell::new((HashMap::new(), 0));
71}
72const DOWNLOADS_FLUSH_MS: u64 = 10_000;
73/// A package's id, and a version's when the download counts for it too.
74type DownloadKey = (String, Option<String>);
75
76thread_local! {
77 /// What billing allows each workspace, as asked last, and when.
78 static ALLOWANCES: RefCell<HashMap<String, (quota::Allowance, u64)>> = RefCell::new(HashMap::new());
79}
80/// How long billing's answer is kept.
81const ALLOWANCE_TTL_MS: u64 = 5 * 60 * 1000;
82
83/// Who made a registry request, as its audit entries and versions name them.
84pub struct Caller {
85 pub actor: Option<AuditActor>,
86}
87
88/// What access decisions need about a package, owned.
89pub struct TargetOf {
90 workspace: String,
91 repo: Option<(String, String, bool)>,
92 public: bool,
93}
94
95impl TargetOf {
96 pub fn package(row: &PackageRow) -> TargetOf {
97 TargetOf {
98 workspace: row.workspace.clone(),
99 repo: row
100 .repo_id
101 .as_ref()
102 .map(|id| (id.clone(), row.repo_name.clone().unwrap_or_default(), !row.public())),
103 public: row.public(),
104 }
105 }
106
107 pub fn view(&self) -> Target<'_> {
108 Target {
109 workspace: &self.workspace,
110 repo: self.repo.as_ref().map(|(id, name, private)| LinkedTo { id, name, private: *private }),
111 public: self.public,
112 }
113 }
114}
115
116pub struct Packages {
117 pub db: Db,
118 pub store: Store,
119 identity: Fetcher,
120 repos: Fetcher,
121 events: Fetcher,
122 /// Signs registry tokens (PACKAGES_TOKEN_SECRET).
123 secret: Vec<u8>,
124 max_request: u64,
125 /// The host package addresses start with: `g1t.sh`.
126 host: String,
127 /// Whether free workspaces' package storage is limited (STORAGE_LIMITS
128 /// `on`); off when self-hosted.
129 storage_limits: bool,
130 env: Env,
131}
132
133/// What a workspace's deletion or restoring does to its packages.
134#[derive(Debug, PartialEq, Eq)]
135pub(crate) enum WorkspaceMark {
136 /// Hide them from `at`, while the workspace may still be restored.
137 Hide { workspace: String, at: String },
138 /// Show them again.
139 Show { workspace: String },
140}
141
142/// The mark `event` asks for, if any. A protected workspace (`protected`,
143/// slugs or ids, as identity's PROTECTED_WORKSPACES) is never hidden: it
144/// cannot be deleted, so an event saying so is a mistake.
145pub(crate) fn workspace_mark(event: &Event, protected: &[String], now: &str) -> Option<WorkspaceMark> {
146 let text = |key: &str| event.data[key].as_str().map(|v| v.trim().to_lowercase()).filter(|v| !v.is_empty());
147 let workspace = text("slug")?;
148 match event.kind.as_str() {
149 "workspace.deleting" => {
150 let id = text("workspaceId").unwrap_or_default();
151 if protected.iter().any(|name| *name == workspace || (!id.is_empty() && *name == id)) {
152 return None;
153 }
154 Some(WorkspaceMark::Hide { workspace, at: now.to_owned() })
155 }
156 "workspace.restored" => Some(WorkspaceMark::Show { workspace }),
157 _ => None,
158 }
159}
160
161fn not_found<T>() -> Outcome<T> {
162 Outcome::fail(FailureCode::NotFound, "Package not found.")
163}
164
165impl Packages {
166 pub fn from_env(env: &Env) -> Result<Packages> {
167 let secret = store::var(env, "PACKAGES_TOKEN_SECRET");
168 if secret.len() < 32 {
169 return Err(worker::Error::RustError("PACKAGES_TOKEN_SECRET is missing or shorter than 32 characters".into()));
170 }
171 // 0 is no limit, as self-hosted (nothing in front cuts bodies short).
172 let max_request = match store::var(env, "MAX_REQUEST_BYTES").parse::<u64>() {
173 Ok(0) => u64::MAX,
174 Ok(limit) => limit,
175 Err(_) => DEFAULT_MAX_REQUEST_BYTES,
176 };
177 let host = store::var(env, "REGISTRY_HOST");
178 Ok(Packages {
179 db: Db { db: env.d1("DB")? },
180 store: store::from_env(env)?,
181 identity: env.service("IDENTITY")?,
182 repos: env.service("REPOS")?,
183 events: env.service("EVENTS")?,
184 secret: secret.into_bytes(),
185 max_request,
186 host: if host.is_empty() { "g1t.sh".to_owned() } else { host },
187 storage_limits: store::var(env, "STORAGE_LIMITS") == "on",
188 env: env.clone(),
189 })
190 }
191
192 async fn viewer_for(&self, username: &str, secret: &str) -> Result<Viewer> {
193 g1t_kit::call(
194 &self.identity,
195 "user_for_git_credentials",
196 &GitCredentialsArgs { username: username.to_owned(), secret: secret.to_owned() },
197 )
198 .await
199 }
200
201 /// A repository of the workspace by name, whoever may see it.
202 async fn repo_by_name(&self, workspace: &str, name: &str) -> Result<Option<Repo>> {
203 if !g1t_contracts::is_valid_repo_name(name) {
204 return Ok(None);
205 }
206 let found: Outcome<Repo> = g1t_kit::call(
207 &self.repos,
208 "get",
209 &GetArgs {
210 path: RepoPath { namespace: workspace.to_owned(), name: name.to_owned() },
211 viewer: Some(User::system(workspace)),
212 },
213 )
214 .await?;
215 Ok(found.into_result().ok().filter(|repo| repo.fork_of.is_none()))
216 }
217
218 /// Who may do what to an image: the package's own settings, or, before
219 /// its first push, the repository it will be linked to.
220 async fn target(&self, name: &names::ImageName, found: Option<&PackageRow>) -> Result<TargetOf> {
221 if let Some(found) = found {
222 return Ok(TargetOf::package(found));
223 }
224 let repo = self.repo_by_name(&name.workspace, name.repo_name()).await?;
225 Ok(TargetOf {
226 workspace: name.workspace.clone(),
227 repo: repo.map(|repo| (repo.id, repo.name, repo.is_private)),
228 public: false,
229 })
230 }
231
232 /// What billing allows the workspace, kept for five minutes unless
233 /// `fresh`, and whether it is the kept answer. `None` when billing
234 /// cannot be asked: the push is then let through.
235 async fn allowance(&self, workspace: &str, fresh: bool) -> Option<(quota::Allowance, bool)> {
236 let now = now_ms();
237 let kept = ALLOWANCES.with(|kept| kept.borrow().get(workspace).copied());
238 if let Some((allowance, at)) = kept
239 && !fresh
240 && now.saturating_sub(at) < ALLOWANCE_TTL_MS
241 {
242 return Some((allowance, true));
243 }
244 let billing = self.env.service("BILLING").ok()?;
245 let asked: Result<quota::Allowance> = g1t_kit::call(
246 &billing,
247 "entitlements",
248 &g1t_contracts::billing::EntitlementsArgs { workspace: workspace.to_owned() },
249 )
250 .await;
251 match asked {
252 Ok(allowance) => {
253 ALLOWANCES.with(|kept| kept.borrow_mut().insert(workspace.to_owned(), (allowance, now)));
254 Some((allowance, false))
255 }
256 Err(error) => {
257 worker::console_error!("packages: billing could not be asked about {workspace}, letting the push through: {error}");
258 None
259 }
260 }
261 }
262
263 /// Why a push of these files (digest and size) into the package may
264 /// not be kept, if it may not: a free workspace past its free storage.
265 /// Files the workspace already holds add nothing.
266 pub(crate) async fn storage_refusal(&self, package: &PackageRow, files: &[(String, u64)]) -> Result<Option<String>> {
267 if !self.storage_limits || files.is_empty() {
268 return Ok(None);
269 }
270 let digests: Vec<String> = files.iter().map(|(digest, _)| digest.clone()).collect();
271 let held = self.db.held(&package.workspace, &digests).await?;
272 let mut seen = std::collections::HashSet::new();
273 let adding: u64 = files
274 .iter()
275 .filter(|(digest, _)| !held.contains(digest) && seen.insert(digest.clone()))
276 .map(|(_, size)| size)
277 .sum();
278 if adding == 0 {
279 return Ok(None);
280 }
281 let Some((allowance, kept)) = self.allowance(&package.workspace, false).await else {
282 return Ok(None);
283 };
284 let (public_bytes, private_bytes) = self.db.storage(&package.workspace).await?;
285 let public = package.public();
286 let used = if public { public_bytes } else { private_bytes };
287 let mut refused = quota::decide(&allowance, public, used, adding).err();
288 // A refusal from the kept answer is checked with billing again: the
289 // workspace may have just added a plan, and must not wait minutes
290 // for the push to go through.
291 if refused.is_some() && kept {
292 refused = match self.allowance(&package.workspace, true).await {
293 Some((allowance, _)) => quota::decide(&allowance, public, used, adding).err(),
294 None => None,
295 };
296 }
297 Ok(refused.map(|refusal| quota::message(&package.workspace, &refusal)))
298 }
299
300 fn count_download(&self, package_id: &str, ctx: &Context) {
301 self.count_downloads(package_id, None, ctx);
302 }
303
304 /// A download of one version, counted for it and its package.
305 fn count_version_download(&self, package_id: &str, version_id: &str, ctx: &Context) {
306 self.count_downloads(package_id, Some(version_id), ctx);
307 }
308
309 fn count_downloads(&self, package_id: &str, version_id: Option<&str>, ctx: &Context) {
310 let due = DOWNLOADS.with(|counts| {
311 let mut counts = counts.borrow_mut();
312 *counts.0.entry((package_id.to_owned(), version_id.map(str::to_owned))).or_default() += 1;
313 let now = now_ms();
314 if now.saturating_sub(counts.1) < DOWNLOADS_FLUSH_MS {
315 return None;
316 }
317 counts.1 = now;
318 Some(counts.0.drain().collect::<Vec<_>>())
319 });
320 if let Some(due) = due {
321 let env = self.env.clone();
322 ctx.wait_until(async move {
323 let written = match env.d1("DB") {
324 Ok(db) => Db { db }.add_downloads(&due).await,
325 Err(error) => Err(error),
326 };
327 if let Err(error) = written {
328 worker::console_error!("packages: downloads not counted: {error}");
329 }
330 });
331 }
332 }
333
334 async fn announce(&self, kind: &'static str, package: &PackageRow, data: PackageEvent, caller: &Caller) {
335 let event = NewEvent {
336 kind,
337 source: SOURCE,
338 repo_id: package.repo_id.clone(),
339 actor: caller.actor.as_ref().map(|actor| actor.actor_id.clone()),
340 data,
341 };
342 let published: Result<serde_json::Value> = g1t_kit::call(&self.events, "publish", &Publish { events: vec![event] }).await;
343 if let Err(error) = published {
344 worker::console_error!("packages: {kind} not published: {error}");
345 }
346 }
347
348 async fn audit(&self, caller: &Caller, action: &str, package: &PackageRow, path: Option<&str>, surface: Option<Surface>) {
349 let Some(actor) = caller.actor.clone() else {
350 return;
351 };
352 let target = AuditTarget {
353 workspace: package.workspace.clone(),
354 repo: package.repo_name.as_ref().map(|name| format!("{}/{name}", package.workspace)),
355 path: Some(path.map_or_else(|| format!("{}:{}/{}", package.ecosystem, package.workspace, package.name), str::to_owned)),
356 ..AuditTarget::default()
357 };
358 let mut entry = NewAuditEntry::new(
359 actor,
360 action,
361 surface.unwrap_or(Surface::Registry),
362 target,
363 &Decision::allow("packages"),
364 new_id("req", now_ms()),
365 );
366 entry.result = Some("ok".to_owned());
367 let recorded: Result<u32> = g1t_kit::call(&self.events, "audit_record", &RecordAuditArgs { entries: vec![entry] }).await;
368 if let Err(error) = recorded {
369 worker::console_error!("packages: audit entry not recorded: {error}");
370 }
371 }
372
373 fn summary(&self, row: &db::ListedRow) -> PackageSummary {
374 let p = &row.package;
375 PackageSummary {
376 id: p.id.clone(),
377 workspace: p.workspace.clone(),
378 ecosystem: Ecosystem::parse(&p.ecosystem).unwrap_or(Ecosystem::Container),
379 name: p.name.clone(),
380 address: match p.ecosystem.as_str() {
381 "npm" => format!("{}/-/npm/@{}/{}", self.host, p.workspace, p.name),
382 "composer" => format!("{}/-/composer/{}/{}", self.host, p.workspace, p.name),
383 "cargo" => format!("{}/-/cargo/{}/{}", self.host, p.workspace, p.name),
384 "maven" => format!("{}/-/maven/{}/{}", self.host, p.workspace, p.name),
385 "nuget" => format!("{}/-/nuget/{}/{}", self.host, p.workspace, p.name),
386 "rubygems" => format!("{}/-/rubygems/{}/{}", self.host, p.workspace, p.name),
387 _ => format!("{}/{}/{}", self.host, p.workspace, p.name),
388 },
389 visibility: Visibility::parse(&p.visibility),
390 repo: p.repo_id.as_ref().map(|id| LinkedRepo {
391 id: id.clone(),
392 namespace: p.workspace.clone(),
393 name: p.repo_name.clone().unwrap_or_default(),
394 }),
395 description: p.description.clone(),
396 versions: row.version_count,
397 latest: db::latest_shown(row),
398 size: row.bytes,
399 downloads: p.downloads,
400 created_at: p.created_at.clone(),
401 updated_at: p.updated_at.clone(),
402 }
403 }
404
405 async fn listed(&self, workspace: &str, ecosystem: Ecosystem, name: &str) -> Result<Option<db::ListedRow>> {
406 let rows = self.db.list(workspace, Some(ecosystem.as_str()), None, None, LIST_LIMIT).await?;
407 if let Some(row) = rows.into_iter().find(|row| row.package.name == name) {
408 return Ok(Some(row));
409 }
410 // Past the first page: read it alone.
411 Ok(self.db.package(workspace, ecosystem.as_str(), name).await?.filter(|package| !package.hidden()).map(|package| db::ListedRow {
412 package,
413 version_count: 0,
414 bytes: 0,
415 latest_tag: None,
416 latest_tag_version: None,
417 latest_version: None,
418 }))
419 }
420
421 async fn list_packages(&self, a: ListPackagesArgs) -> Result<Outcome<Vec<PackageSummary>>> {
422 let workspace = a.workspace.to_lowercase();
423 let rows = self
424 .db
425 .list(&workspace, a.ecosystem.map(Ecosystem::as_str), a.repo_id.as_deref(), a.query.as_deref(), LIST_LIMIT)
426 .await?;
427 let visible = rows
428 .iter()
429 .filter(|row| access::decide(a.viewer.as_ref(), &TargetOf::package(&row.package).view(), Action::Pull).allowed)
430 .map(|row| self.summary(row))
431 .collect();
432 Ok(Outcome::Ok(visible))
433 }
434
435 async fn get_package(&self, a: GetPackageArgs) -> Result<Outcome<PackageDetail>> {
436 let workspace = a.workspace.to_lowercase();
437 let Some(row) = self.listed(&workspace, a.ecosystem, &a.name).await? else {
438 return Ok(not_found());
439 };
440 let target = TargetOf::package(&row.package);
441 let permissions = access::permissions(a.viewer.as_ref(), &target.view());
442 if !permissions.pull {
443 return Ok(not_found());
444 }
445 let tags = self.db.tags(&row.package.id).await?;
446 // A yanked or unlisted version reads as deprecated: still there
447 // for lockfiles that name it, no longer picked for new ones.
448 let withdrawn = match row.package.ecosystem.as_str() {
449 "nuget" => "Unlisted: still restored by projects that name it, no longer shown in search.",
450 "rubygems" => "Yanked: Bundler no longer picks this version for new lockfiles.",
451 _ => "Yanked: Cargo no longer picks this version for new lockfiles.",
452 };
453 let versions = self
454 .db
455 .versions(&row.package.id, VERSIONS_SHOWN)
456 .await?
457 .into_iter()
458 .map(|version| {
459 let meta = version.meta();
460 let text = |key: &str| meta[key].as_str().map(str::to_owned);
461 PackageVersion {
462 tags: tags.iter().filter(|tag| tag.version_id == version.id).map(|tag| tag.tag.clone()).collect(),
463 media_type: text("media_type"),
464 artifact_type: text("artifact_type"),
465 platforms: meta["platforms"]
466 .as_array()
467 .map(|list| list.iter().filter_map(|p| p.as_str().map(str::to_owned)).collect())
468 .unwrap_or_default(),
469 id: version.id,
470 version: version.version,
471 digest: version.digest,
472 size: version.size,
473 subject: version.subject,
474 published_by: version.published_by,
475 published_at: version.published_at,
476 deprecated: if version.yanked != 0 {
477 Some(withdrawn.to_owned())
478 } else {
479 version.deprecated
480 },
481 symbols: meta["symbols"] == true,
482 downloads: (row.package.ecosystem == "nuget").then_some(version.downloads),
483 }
484 })
485 .collect();
486 // The README its page shows: npm's, from the latest version.
487 let readme = match self.db.readme_digest(&row.package.id).await?.and_then(|d| digest::Digest::parse(&d)) {
488 Some(digest) => match self.db.blob(&digest).await? {
489 Some(blob) => self.store.read(&blob.object_key).await?.map(|b| String::from_utf8_lossy(&b).into_owned()),
490 None => None,
491 },
492 None => None,
493 };
494 Ok(Outcome::Ok(PackageDetail {
495 readme,
496 package: self.summary(&row),
497 versions,
498 tags: tags
499 .into_iter()
500 .map(|tag| PackageTag { tag: tag.tag, digest: tag.digest, updated_at: tag.updated_at })
501 .collect(),
502 permissions,
503 }))
504 }
505
506 /// The package `actor` asks to change, if they may `action` it.
507 async fn for_change(&self, actor: &User, workspace: &str, ecosystem: Ecosystem, name: &str, action: Action) -> Result<Outcome<PackageRow>> {
508 let Some(package) = self.db.package(&workspace.to_lowercase(), ecosystem.as_str(), name).await?.filter(|p| !p.hidden()) else {
509 return Ok(not_found());
510 };
511 let target = TargetOf::package(&package);
512 let decision = access::decide(Some(actor), &target.view(), action);
513 if decision.allowed {
514 return Ok(Outcome::Ok(package));
515 }
516 if !access::decide(Some(actor), &target.view(), Action::Pull).allowed {
517 return Ok(not_found());
518 }
519 Ok(Outcome::fail(FailureCode::Forbidden, decision.reason.unwrap_or_else(|| "Not allowed.".to_owned())))
520 }
521
522 async fn delete_version(&self, a: DeleteVersionArgs) -> Result<Outcome<()>> {
523 let package = match self.for_change(&a.actor, &a.workspace, a.ecosystem, &a.name, Action::Delete).await? {
524 Outcome::Ok(package) => package,
525 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
526 };
527 let Some(version) = self.db.find_version(&package.id, &a.version).await? else {
528 return Ok(Outcome::fail(FailureCode::NotFound, "Version not found."));
529 };
530 let caller = Caller { actor: Some(AuditActor::of(&a.actor)) };
531 self.remove_version(&package, &version, &caller).await?;
532 Ok(Outcome::Ok(()))
533 }
534
535 async fn delete_package(&self, a: DeletePackageArgs) -> Result<Outcome<()>> {
536 let package = match self.for_change(&a.actor, &a.workspace, a.ecosystem, &a.name, Action::Delete).await? {
537 Outcome::Ok(package) => package,
538 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
539 };
540 self.db.delete_package(&package.id).await?;
541 self.db.measure(&package.workspace).await?;
542 let caller = Caller { actor: Some(AuditActor::of(&a.actor)) };
543 self.announce("package.deleted", &package, self.event_of(&package), &caller).await;
544 self.audit(&caller, "package.delete", &package, None, a.surface).await;
545 Ok(Outcome::Ok(()))
546 }
547
548 async fn set_package(&self, a: SetPackageArgs) -> Result<Outcome<PackageSummary>> {
549 let package = match self.for_change(&a.actor, &a.workspace, a.ecosystem, &a.name, Action::Delete).await? {
550 Outcome::Ok(package) => package,
551 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
552 };
553 let caller = Caller { actor: Some(AuditActor::of(&a.actor)) };
554 let now = now_ms();
555 let before = package.visibility.clone();
556 if a.unlink {
557 self.db.set_link(&package.id, None, "private", now).await?;
558 } else if let Some(link) = a.link.as_deref() {
559 let Some(repo) = self.repo_by_name(&package.workspace, link).await? else {
560 return Ok(Outcome::fail(FailureCode::NotFound, format!("There is no repository {}/{link}.", package.workspace)));
561 };
562 // Linking hands the package to the repository's roles: only
563 // someone who administers that repository may.
564 let role = g1t_contracts::access::permission(Some(&a.actor), &repo);
565 if role < Some(g1t_contracts::access::RepoRole::Admin) {
566 return Ok(Outcome::fail(
567 FailureCode::Forbidden,
568 format!("You need the Admin role on {}/{} to link a package to it.", repo.namespace, repo.name),
569 ));
570 }
571 let visibility = if repo.is_private { "private" } else { "public" };
572 self.db.set_link(&package.id, Some((&repo.id, &repo.name)), visibility, now).await?;
573 }
574 if let Some(visibility) = a.visibility {
575 let linked = if a.unlink { false } else { a.link.is_some() || package.repo_id.is_some() };
576 if linked {
577 return Ok(Outcome::fail(
578 FailureCode::Invalid,
579 "A linked package has its repository's visibility. Change the repository's, or unlink the package.",
580 ));
581 }
582 self.db.set_visibility(&package.id, visibility.as_str(), now).await?;
583 }
584 self.audit(&caller, "package.update", &package, None, a.surface).await;
585 let Some(row) = self.listed(&package.workspace, a.ecosystem, &package.name).await? else {
586 return Ok(not_found());
587 };
588 if row.package.visibility != before {
589 self.db.measure(&package.workspace).await?;
590 let event = PackageEvent { visibility: Some(row.package.visibility.clone()), ..self.event_of(&row.package) };
591 self.announce("package.visibility_changed", &row.package, event, &caller).await;
592 }
593 Ok(Outcome::Ok(self.summary(&row)))
594 }
595
596 /// Lets go of expired uploads, and of blobs no version has used for a day.
597 async fn sweep(&self) -> Result<(u32, u32)> {
598 let now = now_ms();
599 let mut uploads = 0;
600 for row in self.db.expired_uploads(now, SWEEP_BATCH).await? {
601 if let Some(progress) = row.progress()
602 && let Err(error) = upload::abort(&self.store, &progress).await
603 {
604 worker::console_error!("packages: upload {} not aborted: {error}", row.id);
605 }
606 self.db.delete_upload(&row.id).await?;
607 uploads += 1;
608 }
609 let mut blobs = 0;
610 for blob in self.db.unused_blobs(now, SWEEP_BATCH).await? {
611 self.store.delete(&blob.object_key).await?;
612 self.db.forget_blob(&blob.digest).await?;
613 blobs += 1;
614 }
615 Ok((uploads, blobs))
616 }
617
618 /// Follows what happens elsewhere: workspaces renamed and deleted,
619 /// repositories that change visibility, are renamed, move or go.
620 async fn on_event(&self, env: &Env, event: &Event) -> Result<()> {
621 // Composer packages are read from repositories: again when one is
622 // pushed to (its default branch may have gained a composer.json),
623 // restored, renamed or moved; gone when it is deleted. A failure is
624 // logged, not retried with the batch: the next push reads it again.
625 let repo_id = event.repo_id.clone().or_else(|| event.data["repoId"].as_str().map(str::to_owned));
626 if let Some(repo_id) = repo_id.as_deref() {
627 let synced = match event.kind.as_str() {
628 "git.push" => {
629 let known = self.db.package_for_repo(repo_id, "composer").await?.is_some();
630 if known || event.data["defaultBranch"].as_bool() == Some(true) {
631 Some(self.sync_composer(repo_id).await.map(|_| ()))
632 } else {
633 None
634 }
635 }
636 "repo.deleted" => Some(self.composer_repo_gone(repo_id).await),
637 "repo.restored" | "repo.renamed" | "repo.transferred" => Some(self.sync_composer(repo_id).await.map(|_| ())),
638 _ => None,
639 };
640 if let Some(Err(error)) = synced {
641 worker::console_error!("packages: composer {} for {repo_id}: {error}", event.kind);
642 }
643 if event.kind == "git.push" || event.kind == "repo.deleted" || event.kind == "repo.restored" {
644 return Ok(());
645 }
646 }
647 let db = &self.db.db;
648 let protected = g1t_contracts::identity::protected_names(Some(&store::var(env, "PROTECTED_WORKSPACES")));
649 match workspace_mark(event, &protected, &g1t_contracts::time::rfc3339(now_ms())) {
650 Some(WorkspaceMark::Hide { workspace, at }) => return self.db.mark_workspace(&workspace, Some(&at)).await,
651 Some(WorkspaceMark::Show { workspace }) => return self.db.mark_workspace(&workspace, None).await,
652 None => {}
653 }
654 if g1t_kit::rename::on_event(env, db, event, &[
655 "UPDATE packages SET workspace = ?1 WHERE workspace = ?2",
656 "UPDATE uploads SET workspace = ?1 WHERE workspace = ?2",
657 "UPDATE workspace_blobs SET workspace = ?1 WHERE workspace = ?2",
658 ])
659 .await?
660 {
661 return Ok(());
662 }
663 if g1t_kit::deleted::on_event(db, event, &[
664 "DELETE FROM tags WHERE package_id IN (SELECT id FROM packages WHERE workspace = ?1)",
665 "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)",
666 "DELETE FROM versions WHERE package_id IN (SELECT id FROM packages WHERE workspace = ?1)",
667 "DELETE FROM package_blobs WHERE package_id IN (SELECT id FROM packages WHERE workspace = ?1)",
668 "DELETE FROM uploads WHERE workspace = ?1",
669 "DELETE FROM workspace_blobs WHERE workspace = ?1",
670 "DELETE FROM packages WHERE workspace = ?1",
671 ])
672 .await?
673 {
674 return Ok(());
675 }
676 // A repository renamed keeps its packages linked under its new
677 // name; one moved to another workspace leaves them behind, unlinked.
678 if g1t_kit::transfer::on_event(env, db, event, &[
679 "UPDATE packages SET repo_name = ?6 WHERE repo_id = ?5 AND workspace = ?3",
680 "UPDATE packages SET repo_id = NULL, repo_name = NULL WHERE repo_id = ?5 AND workspace <> ?3",
681 ])
682 .await?
683 {
684 return Ok(());
685 }
686 if g1t_kit::lifecycle::on_purged(db, event, &["UPDATE packages SET repo_id = NULL, repo_name = NULL, visibility = 'private' WHERE repo_id = ?1"])
687 .await?
688 {
689 return Ok(());
690 }
691 if event.kind == "repo.visibility_changed" {
692 #[derive(serde::Deserialize)]
693 #[serde(rename_all = "camelCase")]
694 struct Changed {
695 repo_id: String,
696 is_private: bool,
697 }
698
699 let Ok(changed) = serde_json::from_value::<Changed>(event.data.clone()) else {
700 worker::console_error!("repo.visibility_changed {} could not be read", event.id);
701 return Ok(());
702 };
703 self.db.follow_visibility(&changed.repo_id, changed.is_private).await?;
704 for workspace in self.db.workspaces_linked_to(&changed.repo_id).await? {
705 self.db.measure(&workspace).await?;
706 }
707 }
708 Ok(())
709 }
710}
711
712#[event(fetch)]
713async fn fetch(mut request: Request, env: Env, ctx: Context) -> Result<Response> {
714 let packages = Packages::from_env(&env)?;
715 let Some(method) = rpc_method(&request) else {
716 if request.path().starts_with("/-/npm/") || request.path() == "/-/npm" {
717 return packages.npm(request, &ctx).await;
718 }
719 if request.path().starts_with("/-/composer/") {
720 return packages.composer(request, &ctx).await;
721 }
722 if request.path().starts_with("/-/cargo/") {
723 return packages.cargo(request, &ctx).await;
724 }
725 if request.path().starts_with("/-/maven/") {
726 return packages.maven(request, &ctx).await;
727 }
728 if request.path().starts_with("/-/nuget/") {
729 return packages.nuget(request, &ctx).await;
730 }
731 if request.path().starts_with("/-/rubygems/") {
732 return packages.rubygems(request, &ctx).await;
733 }
734 return packages.registry(request, &ctx).await;
735 };
736 let body: serde_json::Value = request.json().await?;
737 match method.as_str() {
738 "list_packages" => reply(&packages.list_packages(args(body)?).await?),
739 "get_package" => reply(&packages.get_package(args(body)?).await?),
740 "delete_version" => reply(&packages.delete_version(args(body)?).await?),
741 "delete_package" => reply(&packages.delete_package(args(body)?).await?),
742 "set_package" => reply(&packages.set_package(args(body)?).await?),
743 // For billing: what a workspace's packages hold.
744 "storage_all" => reply(&packages.db.storage_all().await?),
745 // Read a repository's Composer package again now, as a push would.
746 "sync_composer" => {
747 let a: SyncComposerArgs = args(body)?;
748 reply(&packages.sync_composer(&a.repo_id).await?)
749 }
750 "storage" => {
751 let a: StorageArgs = args(body)?;
752 let (public_bytes, private_bytes) = packages.db.storage(&a.workspace.to_lowercase()).await?;
753 reply(&PackageStorage { public_bytes, private_bytes })
754 }
755 _ => Response::error("Unknown method", 404),
756 }
757}
758
759/// Every hour: expired uploads are let go, and blobs no version has used
760/// for a day are deleted from the store.
761#[event(scheduled)]
762async fn scheduled(_event: ScheduledEvent, env: Env, _ctx: ScheduleContext) {
763 let packages = match Packages::from_env(&env) {
764 Ok(packages) => packages,
765 Err(error) => {
766 worker::console_error!("packages: the sweep could not start: {error}");
767 return;
768 }
769 };
770 match packages.sweep().await {
771 Ok((0, 0)) => {}
772 Ok((uploads, blobs)) => worker::console_log!("packages: let go of {uploads} uploads and {blobs} blobs"),
773 Err(error) => worker::console_error!("packages: the sweep failed: {error}"),
774 }
775 // Repositories that had a composer.json before the registry did.
776 match packages.composer_backfill().await {
777 Ok(0) => {}
778 Ok(found) => worker::console_log!("packages: found {found} Composer packages"),
779 Err(error) => worker::console_error!("packages: the Composer backfill failed: {error}"),
780 }
781}
782
783#[event(queue)]
784async fn queue(batch: MessageBatch<Event>, env: Env, _ctx: Context) -> Result<()> {
785 let packages = Packages::from_env(&env)?;
786 for message in batch.messages()? {
787 packages.on_event(&env, message.body()).await?;
788 }
789 Ok(())
790}
791
792#[cfg(test)]
793mod tests {
794 use super::*;
795 use serde_json::json;
796
797 fn event(kind: &str, data: serde_json::Value) -> Event {
798 Event {
799 id: "evt_1".into(),
800 kind: kind.into(),
801 source: "identity".into(),
802 time: "2026-10-06T00:00:00.000Z".into(),
803 repo_id: None,
804 actor: None,
805 data,
806 }
807 }
808
809 #[test]
810 fn a_deleted_workspace_is_hidden_and_a_restored_one_shown() {
811 let protected = g1t_contracts::identity::protected_names(Some("wsp_keep"));
812 let now = "2026-10-06T12:00:00.000Z";
813 let deleting = event("workspace.deleting", json!({ "workspaceId": "wsp_1", "slug": "Acme", "by": "ana", "purgeAfter": "x" }));
814 assert_eq!(
815 workspace_mark(&deleting, &protected, now),
816 Some(WorkspaceMark::Hide { workspace: "acme".into(), at: now.into() })
817 );
818 let restored = event("workspace.restored", json!({ "workspaceId": "wsp_1", "slug": "acme" }));
819 assert_eq!(workspace_mark(&restored, &protected, now), Some(WorkspaceMark::Show { workspace: "acme".into() }));
820 }
821
822 #[test]
823 fn protected_workspaces_and_other_events_are_left_alone() {
824 let protected = g1t_contracts::identity::protected_names(Some("wsp_keep"));
825 let now = "2026-10-06T12:00:00.000Z";
826 let flagon = event("workspace.deleting", json!({ "workspaceId": "wsp_f", "slug": "flagon-io" }));
827 assert_eq!(workspace_mark(&flagon, &protected, now), None, "always protected");
828 let by_id = event("workspace.deleting", json!({ "workspaceId": "wsp_keep", "slug": "kept" }));
829 assert_eq!(workspace_mark(&by_id, &protected, now), None, "protected by id");
830 let purged = event("workspace.deleted", json!({ "workspaceId": "wsp_1", "slug": "acme" }));
831 assert_eq!(workspace_mark(&purged, &protected, now), None, "the purge is handled apart");
832 let nameless = event("workspace.deleting", json!({ "workspaceId": "wsp_1", "slug": "" }));
833 assert_eq!(workspace_mark(&nameless, &protected, now), None);
834 }
835}