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