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