| 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 | |
| 10 | mod access; |
| 11 | mod archive; |
| 12 | mod cargo; |
| 13 | mod cargo_http; |
| 14 | mod composer; |
| 15 | mod composer_http; |
| 16 | mod db; |
| 17 | mod digest; |
| 18 | mod limits; |
| 19 | mod manifest; |
| 20 | mod marshal; |
| 21 | mod maven; |
| 22 | mod maven_http; |
| 23 | mod names; |
| 24 | mod npm; |
| 25 | mod npm_http; |
| 26 | mod nuget; |
| 27 | mod nuget_http; |
| 28 | mod oci; |
| 29 | mod quota; |
| 30 | mod range; |
| 31 | mod rubygems; |
| 32 | mod rubygems_http; |
| 33 | mod store; |
| 34 | mod token; |
| 35 | mod upload; |
| 36 | mod xml; |
| 37 | mod yaml; |
| 38 | |
| 39 | use std::cell::RefCell; |
| 40 | use std::collections::HashMap; |
| 41 | |
| 42 | use g1t_contracts::audit::{AuditActor, AuditTarget, NewAuditEntry, RecordAuditArgs, Surface}; |
| 43 | use g1t_contracts::credentials::Decision; |
| 44 | use g1t_contracts::events::{Event, NewEvent, PackageEvent, Publish}; |
| 45 | use g1t_contracts::identity::GitCredentialsArgs; |
| 46 | use g1t_contracts::packages::*; |
| 47 | use g1t_contracts::repos::{GetArgs, Repo, RepoPath}; |
| 48 | use g1t_contracts::{FailureCode, Outcome, User, Viewer, new_id}; |
| 49 | use g1t_kit::{args, now_ms, reply, rpc_method}; |
| 50 | use worker::{Context, Env, Fetcher, MessageBatch, Request, Response, Result, ScheduleContext, ScheduledEvent, event}; |
| 51 | |
| 52 | use access::{Action, LinkedTo, Target}; |
| 53 | use db::{Db, PackageRow}; |
| 54 | use store::{BlobStore, Store}; |
| 55 | |
| 56 | pub(crate) const SOURCE: &str = "packages"; |
| 57 | /// How large a request body may be: Cloudflare's limit on the zone's plan. |
| 58 | const DEFAULT_MAX_REQUEST_BYTES: u64 = 100_000_000; |
| 59 | /// How many packages a listing shows. |
| 60 | const LIST_LIMIT: u32 = 200; |
| 61 | const VERSIONS_SHOWN: u32 = 200; |
| 62 | /// How much one sweep lets go of. |
| 63 | const SWEEP_BATCH: u32 = 200; |
| 64 | |
| 65 | thread_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 | } |
| 72 | const DOWNLOADS_FLUSH_MS: u64 = 10_000; |
| 73 | /// A package's id, and a version's when the download counts for it too. |
| 74 | type DownloadKey = (String, Option<String>); |
| 75 | |
| 76 | thread_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. |
| 81 | const ALLOWANCE_TTL_MS: u64 = 5 * 60 * 1000; |
| 82 | |
| 83 | /// Who made a registry request, as its audit entries and versions name them. |
| 84 | pub struct Caller { |
| 85 | pub actor: Option<AuditActor>, |
| 86 | } |
| 87 | |
| 88 | /// What access decisions need about a package, owned. |
| 89 | pub struct TargetOf { |
| 90 | workspace: String, |
| 91 | repo: Option<(String, String, bool)>, |
| 92 | public: bool, |
| 93 | } |
| 94 | |
| 95 | impl 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 | |
| 116 | pub 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)] |
| 135 | pub(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. |
| 145 | pub(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 | |
| 161 | fn not_found<T>() -> Outcome<T> { |
| 162 | Outcome::fail(FailureCode::NotFound, "Package not found.") |
| 163 | } |
| 164 | |
| 165 | impl 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)] |
| 713 | async 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)] |
| 762 | async 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)] |
| 784 | async 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)] |
| 793 | mod 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 | } |