Skip to content
1,145 linesCodeBlameRaw
1//! The packages service: the registries a workspace publishes to and
2//! installs from, beside its code (docs.g1t.sh/guides/packages/).
3//! Container images first, over OCI Distribution 1.1 on `g1t.sh/v2/`
4//! (oci.rs).
5//!
6//! Other services reach it over `POST /rpc/<method>`; see
7//! `g1t_contracts::packages` for the methods and their arguments. Any
8//! other request is a registry's own protocol. Files are kept through the
9//! `BlobStore` port (store/), metadata in D1 (db.rs).
10
11mod access;
12mod archive;
13mod cargo;
14mod cargo_http;
15mod composer;
16mod composer_http;
17mod db;
18mod digest;
19mod limits;
20mod manifest;
21mod marshal;
22mod maven;
23mod maven_http;
24mod names;
25mod npm;
26mod npm_http;
27mod nuget;
28mod nuget_http;
29mod oci;
30mod quota;
31mod range;
32mod rubygems;
33mod rubygems_http;
34mod store;
35mod token;
36mod upload;
37mod xml;
38mod yaml;
39
40use std::cell::RefCell;
41use std::collections::HashMap;
42
43use g1t_contracts::audit::{AuditActor, AuditTarget, NewAuditEntry, RecordAuditArgs, Surface};
44use g1t_contracts::credentials::Decision;
45use g1t_contracts::events::{Event, NewEvent, PackageEvent, Publish};
46use g1t_contracts::identity::GitCredentialsArgs;
47use g1t_contracts::packages::*;
48use g1t_contracts::repos::{GetArgs, Repo, RepoPath};
49use g1t_contracts::{FailureCode, Outcome, User, Viewer, new_id};
50use g1t_kit::{args, now_ms, reply, rpc_method};
51use worker::{Context, Env, Fetcher, MessageBatch, MessageExt, Request, Response, Result, ScheduleContext, ScheduledEvent, event};
52
53use access::{Action, LinkedTo, Target};
54use db::{Db, PackageRow};
55use store::{BlobStore, Store};
56
57mod settings;
58
59pub(crate) const SOURCE: &str = "packages";
60/// How large a request body may be: Cloudflare's limit on the zone's plan.
61const DEFAULT_MAX_REQUEST_BYTES: u64 = 100_000_000;
62/// How many packages a listing shows.
63const LIST_LIMIT: u32 = 200;
64const VERSIONS_SHOWN: u32 = 200;
65/// How much one sweep lets go of.
66const SWEEP_BATCH: u32 = 200;
67
68thread_local! {
69 /// Pulls counted since the last write, by package (and by version,
70 /// where it is counted too): written at most every few seconds, so a
71 /// busy image costs one write, not one a pull. What an isolate holds
72 /// when it goes away is lost: the count is approximate.
73 static DOWNLOADS: RefCell<(HashMap<DownloadKey, u64>, u64)> = RefCell::new((HashMap::new(), 0));
74}
75const DOWNLOADS_FLUSH_MS: u64 = 10_000;
76/// A package's id, and a version's when the download counts for it too.
77type DownloadKey = (String, Option<String>);
78
79thread_local! {
80 /// What billing allows each workspace, as asked last, and when.
81 static ALLOWANCES: RefCell<HashMap<String, (quota::Allowance, u64)>> = RefCell::new(HashMap::new());
82}
83/// How long billing's answer is kept.
84const ALLOWANCE_TTL_MS: u64 = 5 * 60 * 1000;
85
86/// Who made a registry request, as its audit entries and versions name them.
87pub struct Caller {
88 pub actor: Option<AuditActor>,
89 /// Who it is, where that is known: what linking a package to the
90 /// repository its source label names checks.
91 pub viewer: Option<User>,
92}
93
94impl Caller {
95 pub fn of(viewer: Option<&User>) -> Caller {
96 Caller { actor: viewer.map(AuditActor::of), viewer: viewer.cloned() }
97 }
98
99 pub fn system() -> Caller {
100 Caller { actor: Some(AuditActor::system()), viewer: None }
101 }
102}
103
104/// What access decisions need about a package, owned. Its own grants and
105/// Actions access are read only when a decision needs them (`loaded`).
106pub struct TargetOf {
107 pub(crate) workspace: String,
108 pub(crate) name: String,
109 pub(crate) repo: Option<(String, String, bool)>,
110 pub(crate) public: bool,
111 pub(crate) package_id: Option<String>,
112 pub(crate) inherit: bool,
113 pub(crate) grants: Vec<access::Grant>,
114 pub(crate) actions: Vec<access::RepoAccess>,
115 pub(crate) teams: Vec<String>,
116 pub(crate) loaded: bool,
117}
118
119impl TargetOf {
120 pub fn package(row: &PackageRow) -> TargetOf {
121 TargetOf {
122 workspace: row.workspace.clone(),
123 name: row.name.clone(),
124 repo: row
125 .repo_id
126 .as_ref()
127 .map(|id| (id.clone(), row.repo_name.clone().unwrap_or_default(), !row.public())),
128 public: row.public(),
129 package_id: Some(row.id.clone()),
130 inherit: row.inherits(),
131 grants: Vec::new(),
132 actions: Vec::new(),
133 teams: Vec::new(),
134 loaded: false,
135 }
136 }
137
138 /// A package not made yet, which its first push would link to `repo`.
139 pub fn unmade(workspace: &str, name: &str, repo: Option<(String, String, bool)>) -> TargetOf {
140 TargetOf {
141 workspace: workspace.to_owned(),
142 name: name.to_owned(),
143 repo,
144 public: false,
145 package_id: None,
146 inherit: true,
147 grants: Vec::new(),
148 actions: Vec::new(),
149 teams: Vec::new(),
150 loaded: true,
151 }
152 }
153
154 pub fn view(&self) -> Target<'_> {
155 Target {
156 workspace: &self.workspace,
157 name: &self.name,
158 repo: self.repo.as_ref().map(|(id, name, private)| LinkedTo { id, name, private: *private }),
159 public: self.public,
160 exists: self.package_id.is_some(),
161 inherit: self.inherit,
162 grants: &self.grants,
163 actions: &self.actions,
164 teams: &self.teams,
165 }
166 }
167}
168
169pub struct Packages {
170 pub db: Db,
171 pub store: Store,
172 identity: Fetcher,
173 repos: Fetcher,
174 events: Fetcher,
175 /// Signs registry tokens (PACKAGES_TOKEN_SECRET).
176 secret: Vec<u8>,
177 max_request: u64,
178 /// The host package addresses start with: `g1t.sh`.
179 host: String,
180 /// Whether free workspaces' package storage is limited (STORAGE_LIMITS
181 /// `on`); off when self-hosted.
182 storage_limits: bool,
183 env: Env,
184}
185
186/// What a workspace's deletion or restoring does to its packages.
187#[derive(Debug, PartialEq, Eq)]
188pub(crate) enum WorkspaceMark {
189 /// Hide them from `at`, while the workspace may still be restored.
190 Hide { workspace: String, at: String },
191 /// Show them again.
192 Show { workspace: String },
193}
194
195/// The mark `event` asks for, if any. A protected workspace (`protected`,
196/// slugs or ids, as identity's PROTECTED_WORKSPACES) is never hidden: it
197/// cannot be deleted, so an event saying so is a mistake.
198pub(crate) fn workspace_mark(event: &Event, protected: &[String], now: &str) -> Option<WorkspaceMark> {
199 let text = |key: &str| event.data[key].as_str().map(|v| v.trim().to_lowercase()).filter(|v| !v.is_empty());
200 let workspace = text("slug")?;
201 match event.kind.as_str() {
202 "workspace.deleting" => {
203 let id = text("workspaceId").unwrap_or_default();
204 if protected.iter().any(|name| *name == workspace || (!id.is_empty() && *name == id)) {
205 return None;
206 }
207 Some(WorkspaceMark::Hide { workspace, at: now.to_owned() })
208 }
209 "workspace.restored" => Some(WorkspaceMark::Show { workspace }),
210 _ => None,
211 }
212}
213
214/// Adds one download to the counts kept since the last write, and answers
215/// with all of them to write when the last write was long enough ago.
216pub(crate) fn tally(counts: &mut (HashMap<DownloadKey, u64>, u64), key: DownloadKey, now: u64) -> Option<Vec<(DownloadKey, u64)>> {
217 *counts.0.entry(key).or_default() += 1;
218 if now.saturating_sub(counts.1) < DOWNLOADS_FLUSH_MS {
219 return None;
220 }
221 counts.1 = now;
222 Some(counts.0.drain().collect())
223}
224
225fn not_found<T>() -> Outcome<T> {
226 Outcome::fail(FailureCode::NotFound, "Package not found.")
227}
228
229impl Packages {
230 pub fn from_env(env: &Env) -> Result<Packages> {
231 let secret = store::var(env, "PACKAGES_TOKEN_SECRET");
232 if secret.len() < 32 {
233 return Err(worker::Error::RustError("PACKAGES_TOKEN_SECRET is missing or shorter than 32 characters".into()));
234 }
235 // 0 is no limit, as self-hosted (nothing in front cuts bodies short).
236 let max_request = match store::var(env, "MAX_REQUEST_BYTES").parse::<u64>() {
237 Ok(0) => u64::MAX,
238 Ok(limit) => limit,
239 Err(_) => DEFAULT_MAX_REQUEST_BYTES,
240 };
241 let host = store::var(env, "REGISTRY_HOST");
242 Ok(Packages {
243 db: Db { db: env.d1("DB")? },
244 store: store::from_env(env)?,
245 identity: env.service("IDENTITY")?,
246 repos: env.service("REPOS")?,
247 events: env.service("EVENTS")?,
248 secret: secret.into_bytes(),
249 max_request,
250 host: if host.is_empty() { "g1t.sh".to_owned() } else { host },
251 storage_limits: store::var(env, "STORAGE_LIMITS") == "on",
252 env: env.clone(),
253 })
254 }
255
256 async fn viewer_for(&self, username: &str, secret: &str) -> Result<Viewer> {
257 let viewer: Viewer = g1t_kit::call(
258 &self.identity,
259 "user_for_git_credentials",
260 &GitCredentialsArgs { username: username.to_owned(), secret: secret.to_owned() },
261 )
262 .await?;
263 // An account that has not confirmed its email address signs in to
264 // nothing yet: its credentials are refused like wrong ones.
265 Ok(viewer.filter(|user| !user.awaits_confirmation()))
266 }
267
268 /// A repository of the workspace by name, whoever may see it.
269 async fn repo_by_name(&self, workspace: &str, name: &str) -> Result<Option<Repo>> {
270 if !g1t_contracts::is_valid_repo_name(name) {
271 return Ok(None);
272 }
273 let found: Outcome<Repo> = g1t_kit::call(
274 &self.repos,
275 "get",
276 &GetArgs {
277 path: RepoPath { namespace: workspace.to_owned(), name: name.to_owned() },
278 viewer: Some(User::system(workspace)),
279 },
280 )
281 .await?;
282 Ok(found.into_result().ok().filter(|repo| repo.fork_of.is_none()))
283 }
284
285 /// Who may do what to an image: the package's own settings, or, before
286 /// its first push, the repository it will be linked to.
287 async fn target(&self, name: &names::ImageName, found: Option<&PackageRow>) -> Result<TargetOf> {
288 if let Some(found) = found {
289 return Ok(TargetOf::package(found));
290 }
291 let repo = self.repo_by_name(&name.workspace, name.repo_name()).await?;
292 Ok(TargetOf::unmade(&name.workspace, &name.name, repo.map(|repo| (repo.id, repo.name, repo.is_private))))
293 }
294
295 /// What billing allows the workspace, kept for five minutes unless
296 /// `fresh`, and whether it is the kept answer. `None` when billing
297 /// cannot be asked: the push is then let through.
298 async fn allowance(&self, workspace: &str, fresh: bool) -> Option<(quota::Allowance, bool)> {
299 let now = now_ms();
300 let kept = ALLOWANCES.with(|kept| kept.borrow().get(workspace).copied());
301 if let Some((allowance, at)) = kept
302 && !fresh
303 && now.saturating_sub(at) < ALLOWANCE_TTL_MS
304 {
305 return Some((allowance, true));
306 }
307 let billing = self.env.service("BILLING").ok()?;
308 let asked: Result<quota::Allowance> = g1t_kit::call(
309 &billing,
310 "entitlements",
311 &g1t_contracts::billing::EntitlementsArgs { workspace: workspace.to_owned() },
312 )
313 .await;
314 match asked {
315 Ok(allowance) => {
316 ALLOWANCES.with(|kept| kept.borrow_mut().insert(workspace.to_owned(), (allowance, now)));
317 Some((allowance, false))
318 }
319 Err(error) => {
320 worker::console_error!("packages: billing could not be asked about {workspace}, letting the push through: {error}");
321 None
322 }
323 }
324 }
325
326 /// Why a push of these files (digest and size) into the package may
327 /// not be kept, if it may not: a free workspace past its free storage.
328 /// Files the workspace already holds add nothing.
329 pub(crate) async fn storage_refusal(&self, package: &PackageRow, files: &[(String, u64)]) -> Result<Option<String>> {
330 if !self.storage_limits || files.is_empty() {
331 return Ok(None);
332 }
333 let digests: Vec<String> = files.iter().map(|(digest, _)| digest.clone()).collect();
334 let held = self.db.held(&package.workspace, &digests).await?;
335 let mut seen = std::collections::HashSet::new();
336 let adding: u64 = files
337 .iter()
338 .filter(|(digest, _)| !held.contains(digest) && seen.insert(digest.clone()))
339 .map(|(_, size)| size)
340 .sum();
341 if adding == 0 {
342 return Ok(None);
343 }
344 let Some((allowance, kept)) = self.allowance(&package.workspace, false).await else {
345 return Ok(None);
346 };
347 let (public_bytes, private_bytes) = self.db.storage(&package.workspace).await?;
348 let public = package.public();
349 let used = if public { public_bytes } else { private_bytes };
350 let mut refused = quota::decide(&allowance, public, used, adding).err();
351 // A refusal from the kept answer is checked with billing again: the
352 // workspace may have just added a plan, and must not wait minutes
353 // for the push to go through.
354 if refused.is_some() && kept {
355 refused = match self.allowance(&package.workspace, true).await {
356 Some((allowance, _)) => quota::decide(&allowance, public, used, adding).err(),
357 None => None,
358 };
359 }
360 Ok(refused.map(|refusal| quota::message(&package.workspace, &refusal)))
361 }
362
363 fn count_download(&self, package_id: &str, ctx: &Context) {
364 self.count_downloads(package_id, None, ctx);
365 }
366
367 /// A download of one version, counted for it and its package.
368 fn count_version_download(&self, package_id: &str, version_id: &str, ctx: &Context) {
369 self.count_downloads(package_id, Some(version_id), ctx);
370 }
371
372 fn count_downloads(&self, package_id: &str, version_id: Option<&str>, ctx: &Context) {
373 let key = (package_id.to_owned(), version_id.map(str::to_owned));
374 let due = DOWNLOADS.with(|counts| tally(&mut counts.borrow_mut(), key, now_ms()));
375 if let Some(due) = due {
376 let env = self.env.clone();
377 ctx.wait_until(async move {
378 let written = match env.d1("DB") {
379 Ok(db) => Db { db }.add_downloads(&due).await,
380 Err(error) => Err(error),
381 };
382 if let Err(error) = written {
383 worker::console_error!("packages: downloads not counted: {error}");
384 }
385 });
386 }
387 }
388
389 async fn announce(&self, kind: &'static str, package: &PackageRow, data: PackageEvent, caller: &Caller) {
390 let event = NewEvent {
391 kind,
392 source: SOURCE,
393 repo_id: package.repo_id.clone(),
394 actor: caller.actor.as_ref().map(|actor| actor.actor_id.clone()),
395 data,
396 };
397 let published: Result<serde_json::Value> = g1t_kit::call(&self.events, "publish", &Publish { events: vec![event] }).await;
398 if let Err(error) = published {
399 worker::console_error!("packages: {kind} not published: {error}");
400 }
401 }
402
403 async fn audit(&self, caller: &Caller, action: &str, package: &PackageRow, path: Option<&str>, surface: Option<Surface>) {
404 self.audit_with(caller, action, package, path, surface, None).await;
405 }
406
407 /// An audit entry, with what changed in words (`message`), as identity
408 /// writes its access entries.
409 async fn audit_with(
410 &self,
411 caller: &Caller,
412 action: &str,
413 package: &PackageRow,
414 path: Option<&str>,
415 surface: Option<Surface>,
416 message: Option<String>,
417 ) {
418 let Some(actor) = caller.actor.clone() else {
419 return;
420 };
421 let target = AuditTarget {
422 workspace: package.workspace.clone(),
423 repo: package.repo_name.as_ref().map(|name| format!("{}/{name}", package.workspace)),
424 path: Some(path.map_or_else(|| format!("{}:{}/{}", package.ecosystem, package.workspace, package.name), str::to_owned)),
425 ..AuditTarget::default()
426 };
427 let mut entry = NewAuditEntry::new(
428 actor,
429 action,
430 surface.unwrap_or(Surface::Registry),
431 target,
432 &Decision::allow("packages"),
433 new_id("req", now_ms()),
434 );
435 entry.result = Some("ok".to_owned());
436 entry.message = message;
437 let recorded: Result<u32> = g1t_kit::call(&self.events, "audit_record", &RecordAuditArgs { entries: vec![entry] }).await;
438 if let Err(error) = recorded {
439 worker::console_error!("packages: audit entry not recorded: {error}");
440 }
441 }
442
443 fn summary(&self, row: &db::ListedRow) -> PackageSummary {
444 let p = &row.package;
445 PackageSummary {
446 id: p.id.clone(),
447 workspace: p.workspace.clone(),
448 ecosystem: Ecosystem::parse(&p.ecosystem).unwrap_or(Ecosystem::Container),
449 name: p.name.clone(),
450 address: match p.ecosystem.as_str() {
451 "npm" => format!("{}/-/npm/@{}/{}", self.host, p.workspace, p.name),
452 "composer" => format!("{}/-/composer/{}/{}", self.host, p.workspace, p.name),
453 "cargo" => format!("{}/-/cargo/{}/{}", self.host, p.workspace, p.name),
454 "maven" => format!("{}/-/maven/{}/{}", self.host, p.workspace, p.name),
455 "nuget" => format!("{}/-/nuget/{}/{}", self.host, p.workspace, p.name),
456 "rubygems" => format!("{}/-/rubygems/{}/{}", self.host, p.workspace, p.name),
457 _ => format!("{}/{}/{}", self.host, p.workspace, p.name),
458 },
459 visibility: Visibility::parse(&p.visibility),
460 repo: p.repo_id.as_ref().map(|id| LinkedRepo {
461 id: id.clone(),
462 namespace: p.workspace.clone(),
463 name: p.repo_name.clone().unwrap_or_default(),
464 }),
465 description: p.description.clone(),
466 versions: row.version_count,
467 latest: db::latest_shown(row),
468 size: row.bytes,
469 downloads: p.downloads,
470 created_at: p.created_at.clone(),
471 updated_at: p.updated_at.clone(),
472 inherit_access: p.inherits(),
473 deleted_at: p.deleted_at.clone(),
474 deleted_by: p.deleted_by.clone(),
475 purge_at: p.deleted_at.as_deref().and_then(settings::purge_at),
476 }
477 }
478
479 async fn listed(&self, workspace: &str, ecosystem: Ecosystem, name: &str) -> Result<Option<db::ListedRow>> {
480 let rows = self.db.list(workspace, Some(ecosystem.as_str()), None, None, LIST_LIMIT).await?;
481 if let Some(row) = rows.into_iter().find(|row| row.package.name == name) {
482 return Ok(Some(row));
483 }
484 // Past the first page: read it alone.
485 Ok(self.db.package(workspace, ecosystem.as_str(), name).await?.filter(|package| !package.hidden()).map(|package| db::ListedRow {
486 package,
487 version_count: 0,
488 bytes: 0,
489 latest_tag: None,
490 latest_tag_version: None,
491 latest_version: None,
492 }))
493 }
494
495 async fn list_packages(&self, a: ListPackagesArgs) -> Result<Outcome<Vec<PackageSummary>>> {
496 let workspace = a.workspace.to_lowercase();
497 let rows = self
498 .db
499 .list(&workspace, a.ecosystem.map(Ecosystem::as_str), a.repo_id.as_deref(), a.query.as_deref(), LIST_LIMIT)
500 .await?;
501 let packages: Vec<&PackageRow> = rows.iter().map(|row| &row.package).collect();
502 let may = self.may_all(a.viewer.as_ref(), &packages, Action::Pull).await?;
503 let visible = rows
504 .iter()
505 .zip(may)
506 .filter(|(_, may)| *may)
507 .map(|(row, _)| self.summary(row))
508 .collect();
509 Ok(Outcome::Ok(visible))
510 }
511
512 async fn get_package(&self, a: GetPackageArgs) -> Result<Outcome<PackageDetail>> {
513 let workspace = a.workspace.to_lowercase();
514 let Some(row) = self.listed(&workspace, a.ecosystem, &a.name).await? else {
515 return Ok(not_found());
516 };
517 let mut target = TargetOf::package(&row.package);
518 let permissions = self.permissions(a.viewer.as_ref(), &mut target).await?;
519 if !permissions.pull {
520 return Ok(not_found());
521 }
522 let tags = self.db.tags(&row.package.id).await?;
523 let versions = self
524 .db
525 .versions(&row.package.id, VERSIONS_SHOWN)
526 .await?
527 .into_iter()
528 .map(|version| settings::version_of(&row.package, version, &tags))
529 .collect();
530 // The README its page shows: npm's, from the latest version.
531 let readme = match self.db.readme_digest(&row.package.id).await?.and_then(|d| digest::Digest::parse(&d)) {
532 Some(digest) => match self.db.blob(&digest).await? {
533 Some(blob) => self.store.read(&blob.object_key).await?.map(|b| String::from_utf8_lossy(&b).into_owned()),
534 None => None,
535 },
536 None => None,
537 };
538 Ok(Outcome::Ok(PackageDetail {
539 readme,
540 package: self.summary(&row),
541 versions,
542 tags: tags
543 .into_iter()
544 .map(|tag| PackageTag { tag: tag.tag, digest: tag.digest, updated_at: tag.updated_at })
545 .collect(),
546 permissions,
547 }))
548 }
549
550 /// The package `actor` asks to change, if they may `action` it.
551 pub(crate) async fn for_change(&self, actor: &User, workspace: &str, ecosystem: Ecosystem, name: &str, action: Action) -> Result<Outcome<PackageRow>> {
552 let Some(package) = self.db.package(&workspace.to_lowercase(), ecosystem.as_str(), name).await?.filter(|p| !p.hidden()) else {
553 return Ok(not_found());
554 };
555 match self.allowed(actor, &package, action).await? {
556 Outcome::Ok(()) => Ok(Outcome::Ok(package)),
557 Outcome::Fail(failure) => Ok(Outcome::Fail(failure)),
558 }
559 }
560
561 /// Whether `actor` may `action` the package: not found for one who may
562 /// not pull it, forbidden with the reason for one who may.
563 pub(crate) async fn allowed(&self, actor: &User, package: &PackageRow, action: Action) -> Result<Outcome<()>> {
564 let mut target = TargetOf::package(package);
565 let decision = self.decide(Some(actor), &mut target, action).await?;
566 if decision.allowed {
567 return Ok(Outcome::Ok(()));
568 }
569 if !self.decide(Some(actor), &mut target, Action::Pull).await?.allowed {
570 return Ok(not_found());
571 }
572 Ok(Outcome::fail(FailureCode::Forbidden, decision.reason.unwrap_or_else(|| "Not allowed.".to_owned())))
573 }
574
575 async fn delete_version(&self, a: DeleteVersionArgs) -> Result<Outcome<()>> {
576 let package = match self.for_change(&a.actor, &a.workspace, a.ecosystem, &a.name, Action::Delete).await? {
577 Outcome::Ok(package) => package,
578 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
579 };
580 if package.ecosystem == "composer" {
581 return Ok(Outcome::fail(
582 FailureCode::Invalid,
583 "A Composer package's versions are its repository's tags and branches: delete the tag or branch instead.",
584 ));
585 }
586 let Some(version) = self.db.find_version(&package.id, &a.version).await? else {
587 return Ok(Outcome::fail(FailureCode::NotFound, "Version not found."));
588 };
589 let caller = Caller::of(Some(&a.actor));
590 self.remove_version_from(&package, &version, &caller, a.surface).await?;
591 Ok(Outcome::Ok(()))
592 }
593
594 async fn delete_package(&self, a: DeletePackageArgs) -> Result<Outcome<()>> {
595 let package = match self.for_change(&a.actor, &a.workspace, a.ecosystem, &a.name, Action::Delete).await? {
596 Outcome::Ok(package) => package,
597 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
598 };
599 let caller = Caller::of(Some(&a.actor));
600 self.remove_package(&package, &caller, a.surface).await?;
601 Ok(Outcome::Ok(()))
602 }
603
604 /// Deletes a package: hidden at once, restorable for 30 days, its name
605 /// kept until then. Its files go with the purge.
606 pub(crate) async fn remove_package(&self, package: &PackageRow, caller: &Caller, surface: Option<Surface>) -> Result<()> {
607 let by = caller.actor.as_ref().map_or("", |actor| actor.actor.as_str());
608 self.db.soft_delete_package(&package.id, by, now_ms()).await?;
609 self.db.measure(&package.workspace).await?;
610 self.announce("package.deleted", package, self.event_of(package), caller).await;
611 self.audit_with(
612 caller,
613 "package.delete",
614 package,
615 None,
616 surface,
617 Some(format!("Deleted the package; it can be restored for {RESTORE_DAYS} days")),
618 )
619 .await;
620 Ok(())
621 }
622
623 async fn set_package(&self, a: SetPackageArgs) -> Result<Outcome<PackageSummary>> {
624 let package = match self.for_change(&a.actor, &a.workspace, a.ecosystem, &a.name, Action::Admin).await? {
625 Outcome::Ok(package) => package,
626 Outcome::Fail(failure) => return Ok(Outcome::Fail(failure)),
627 };
628 let caller = Caller::of(Some(&a.actor));
629 let now = now_ms();
630 let before = package.visibility.clone();
631 let linked = if a.unlink { false } else { a.link.is_some() || package.repo_id.is_some() };
632 if a.visibility.is_some_and(|visibility| visibility.as_str() != package.visibility) && linked {
633 return Ok(Outcome::fail(
634 FailureCode::Invalid,
635 "A linked package has its repository's visibility. Change the repository's, or unlink the package.",
636 ));
637 }
638 if a.inherit_access.is_some() && !linked {
639 return Ok(Outcome::fail(
640 FailureCode::Invalid,
641 "Only a package linked to a repository inherits access. Link it to a repository first.",
642 ));
643 }
644 let path = |repo: &str| format!("{}/{repo}", package.workspace);
645 if a.unlink {
646 if let Some(name) = package.repo_name.as_deref() {
647 self.db.set_link(&package.id, None, "private", now).await?;
648 self.audit_with(&caller, "package.unlinked", &package, None, a.surface, Some(format!("Unlinked from {}", path(name)))).await;
649 }
650 } else if let Some(link) = a.link.as_deref() {
651 let link = link.trim().trim_end_matches(".git");
652 let name = match link.split_once('/') {
653 Some((owner, name)) if owner.eq_ignore_ascii_case(&package.workspace) => name,
654 Some(_) => {
655 return Ok(Outcome::fail(
656 FailureCode::Invalid,
657 format!("A package can only be linked to a repository of its own workspace, {}.", package.workspace),
658 ));
659 }
660 None => link,
661 };
662 let Some(repo) = self.repo_by_name(&package.workspace, name).await? else {
663 return Ok(Outcome::fail(FailureCode::NotFound, format!("There is no repository {}.", path(name))));
664 };
665 // Linking hands the package to the repository's roles: only
666 // someone who administers that repository may.
667 let role = g1t_contracts::access::permission(Some(&a.actor), &repo);
668 if role < Some(g1t_contracts::access::RepoRole::Admin) {
669 return Ok(Outcome::fail(
670 FailureCode::Forbidden,
671 format!("You need the Admin role on {}/{} to link a package to it.", repo.namespace, repo.name),
672 ));
673 }
674 if package.repo_id.as_deref() != Some(repo.id.as_str()) {
675 let visibility = if repo.is_private { "private" } else { "public" };
676 self.db.set_link(&package.id, Some((&repo.id, &repo.name)), visibility, now).await?;
677 self.audit_with(&caller, "package.linked", &package, None, a.surface, Some(format!("Linked to {}", path(&repo.name)))).await;
678 }
679 }
680 if let Some(visibility) = a.visibility
681 && !linked
682 && visibility.as_str() != package.visibility
683 {
684 self.db.set_visibility(&package.id, visibility.as_str(), now).await?;
685 }
686 if let Some(inherit) = a.inherit_access
687 && inherit != package.inherits()
688 {
689 self.db.set_inherit(&package.id, inherit, now).await?;
690 let message = if inherit {
691 "Turned on inheriting access from the linked repository"
692 } else {
693 "Turned off inheriting access from the linked repository"
694 };
695 self.audit_with(&caller, "package.inherit_access_changed", &package, None, a.surface, Some(message.to_owned())).await;
696 }
697 let Some(row) = self.listed(&package.workspace, a.ecosystem, &package.name).await? else {
698 return Ok(not_found());
699 };
700 if row.package.visibility != before {
701 self.db.measure(&package.workspace).await?;
702 let event = PackageEvent { visibility: Some(row.package.visibility.clone()), ..self.event_of(&row.package) };
703 self.announce("package.visibility_changed", &row.package, event, &caller).await;
704 self.audit_with(
705 &caller,
706 "package.visibility_changed",
707 &row.package,
708 None,
709 a.surface,
710 Some(format!("Made the package {}", row.package.visibility)),
711 )
712 .await;
713 }
714 Ok(Outcome::Ok(self.summary(&row)))
715 }
716
717 /// Lets go of expired uploads, and of blobs no version has used for a day.
718 async fn sweep(&self) -> Result<(u32, u32)> {
719 let now = now_ms();
720 let mut uploads = 0;
721 for row in self.db.expired_uploads(now, SWEEP_BATCH).await? {
722 if let Some(progress) = row.progress()
723 && let Err(error) = upload::abort(&self.store, &progress).await
724 {
725 worker::console_error!("packages: upload {} not aborted: {error}", row.id);
726 }
727 self.db.delete_upload(&row.id).await?;
728 uploads += 1;
729 }
730 let mut blobs = 0;
731 for blob in self.db.unused_blobs(now, SWEEP_BATCH).await? {
732 self.store.delete(&blob.object_key).await?;
733 self.db.forget_blob(&blob.digest).await?;
734 blobs += 1;
735 }
736 Ok((uploads, blobs))
737 }
738
739 /// Follows what happens elsewhere: workspaces renamed and deleted,
740 /// repositories that change visibility, are renamed, move or go.
741 async fn on_event(&self, env: &Env, event: &Event) -> Result<()> {
742 // Composer packages are read from repositories: again when one is
743 // pushed to (its default branch may have gained a composer.json),
744 // restored, renamed or moved; gone when it is deleted. A failure is
745 // logged, not retried with the batch: the next push reads it again.
746 let repo_id = event.repo_id.clone().or_else(|| event.data["repoId"].as_str().map(str::to_owned));
747 if let Some(repo_id) = repo_id.as_deref() {
748 let synced = match event.kind.as_str() {
749 "git.push" => {
750 let known = self.db.package_for_repo(repo_id, "composer").await?.is_some();
751 if known || event.data["defaultBranch"].as_bool() == Some(true) {
752 Some(self.sync_composer(repo_id).await.map(|_| ()))
753 } else {
754 None
755 }
756 }
757 "repo.deleted" => Some(self.composer_repo_gone(repo_id).await),
758 "repo.restored" | "repo.renamed" | "repo.transferred" => Some(self.sync_composer(repo_id).await.map(|_| ())),
759 _ => None,
760 };
761 if let Some(Err(error)) = synced {
762 worker::console_error!("packages: composer {} for {repo_id}: {error}", event.kind);
763 }
764 if event.kind == "git.push" || event.kind == "repo.deleted" || event.kind == "repo.restored" {
765 return Ok(());
766 }
767 }
768 let db = &self.db.db;
769 let protected = g1t_contracts::identity::protected_names(Some(&store::var(env, "PROTECTED_WORKSPACES")));
770 match workspace_mark(event, &protected, &g1t_contracts::time::rfc3339(now_ms())) {
771 Some(WorkspaceMark::Hide { workspace, at }) => return self.db.mark_workspace(&workspace, Some(&at)).await,
772 Some(WorkspaceMark::Show { workspace }) => return self.db.mark_workspace(&workspace, None).await,
773 None => {}
774 }
775 if g1t_kit::rename::on_event(env, db, event, &[
776 "UPDATE packages SET workspace = ?1 WHERE workspace = ?2",
777 "UPDATE uploads SET workspace = ?1 WHERE workspace = ?2",
778 "UPDATE workspace_blobs SET workspace = ?1 WHERE workspace = ?2",
779 ])
780 .await?
781 {
782 return Ok(());
783 }
784 if g1t_kit::deleted::on_event(db, event, &[
785 "DELETE FROM tags WHERE package_id IN (SELECT id FROM packages WHERE workspace = ?1)",
786 "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)",
787 "DELETE FROM versions WHERE package_id IN (SELECT id FROM packages WHERE workspace = ?1)",
788 "DELETE FROM package_blobs WHERE package_id IN (SELECT id FROM packages WHERE workspace = ?1)",
789 "DELETE FROM package_access WHERE package_id IN (SELECT id FROM packages WHERE workspace = ?1)",
790 "DELETE FROM package_actions_access WHERE package_id IN (SELECT id FROM packages WHERE workspace = ?1)",
791 "DELETE FROM uploads WHERE workspace = ?1",
792 "DELETE FROM workspace_blobs WHERE workspace = ?1",
793 "DELETE FROM packages WHERE workspace = ?1",
794 ])
795 .await?
796 {
797 return Ok(());
798 }
799 // A repository renamed keeps its packages linked under its new
800 // name; one moved to another workspace leaves them behind, unlinked.
801 // Manage Actions access follows a rename, and lets go of a
802 // repository that left the package's workspace.
803 if g1t_kit::transfer::on_event(env, db, event, &[
804 "UPDATE packages SET repo_name = ?6 WHERE repo_id = ?5 AND workspace = ?3",
805 "UPDATE packages SET repo_id = NULL, repo_name = NULL WHERE repo_id = ?5 AND workspace <> ?3",
806 "UPDATE package_actions_access SET repo_name = ?6 WHERE repo_id = ?5",
807 "DELETE FROM package_actions_access WHERE repo_id = ?5 AND package_id IN (SELECT id FROM packages WHERE workspace <> ?3)",
808 ])
809 .await?
810 {
811 return Ok(());
812 }
813 if g1t_kit::lifecycle::on_purged(db, event, &[
814 "UPDATE packages SET repo_id = NULL, repo_name = NULL, visibility = 'private' WHERE repo_id = ?1",
815 "DELETE FROM package_actions_access WHERE repo_id = ?1",
816 ])
817 .await?
818 {
819 return Ok(());
820 }
821 // A team renamed or deleted: its roles on packages follow.
822 if matches!(event.kind.as_str(), "team.edited" | "team.deleted") {
823 let field = |snake: &str, camel: &str| {
824 event.data[snake].as_str().or_else(|| event.data[camel].as_str()).map(str::to_owned)
825 };
826 if let Some(team_id) = field("team_id", "teamId") {
827 let slug = field("team", "team");
828 let renamed = (event.kind == "team.edited").then_some(slug.as_deref()).flatten();
829 if event.kind == "team.deleted" || renamed.is_some() {
830 self.db.follow_team(&team_id, renamed).await?;
831 }
832 }
833 return Ok(());
834 }
835 if event.kind == "repo.visibility_changed" {
836 #[derive(serde::Deserialize)]
837 #[serde(rename_all = "camelCase")]
838 struct Changed {
839 repo_id: String,
840 is_private: bool,
841 }
842
843 let Ok(changed) = serde_json::from_value::<Changed>(event.data.clone()) else {
844 worker::console_error!("repo.visibility_changed {} could not be read", event.id);
845 return Ok(());
846 };
847 self.db.follow_visibility(&changed.repo_id, changed.is_private).await?;
848 for workspace in self.db.workspaces_linked_to(&changed.repo_id).await? {
849 self.db.measure(&workspace).await?;
850 }
851 }
852 Ok(())
853 }
854}
855
856impl Packages {
857 /// A registry request, to the registry its path names.
858 async fn serve_registry(&self, request: Request, ctx: &Context) -> Result<Response> {
859 let path = request.path();
860 if path.starts_with("/-/npm/") || path == "/-/npm" {
861 return self.npm(request, ctx).await;
862 }
863 if path.starts_with("/-/composer/") {
864 return self.composer(request, ctx).await;
865 }
866 if path.starts_with("/-/cargo/") {
867 return self.cargo(request, ctx).await;
868 }
869 if path.starts_with("/-/maven/") {
870 return self.maven(request, ctx).await;
871 }
872 if path.starts_with("/-/nuget/") {
873 return self.nuget(request, ctx).await;
874 }
875 if path.starts_with("/-/rubygems/") {
876 return self.rubygems(request, ctx).await;
877 }
878 self.registry(request, ctx).await
879 }
880}
881
882/// The policy for registry answers: what they serve is a publisher's bytes,
883/// on the site's origin, so nothing in them may load or run.
884const NOTHING_RUNS: &str = "default-src 'none'; sandbox";
885
886/// Whether a browser could open a body of this type as a page: HTML, SVG,
887/// any XML, or a type it does not know as data. Such a body is a download.
888fn opens_as_document(content_type: Option<&str>) -> bool {
889 let kind = content_type.unwrap_or("").split(';').next().unwrap_or("").trim().to_ascii_lowercase();
890 if kind.ends_with("+xml") || kind.ends_with("/xml") || kind.contains("html") || kind.contains("svg") || kind.contains("xsl") {
891 return true;
892 }
893 let data = kind == "text/plain"
894 || kind == "application/json"
895 || (kind.starts_with("application/") && kind.ends_with("+json"))
896 || matches!(
897 kind.as_str(),
898 "application/octet-stream"
899 | "application/gzip"
900 | "application/x-gzip"
901 | "application/zip"
902 | "application/x-tar"
903 | "application/java-archive"
904 | "application/pgp-signature"
905 | "image/png"
906 | "image/jpeg"
907 | "image/gif"
908 | "image/webp"
909 | "image/avif"
910 )
911 || kind.starts_with("application/vnd.");
912 !data
913}
914
915/// The headers every registry answer carries: no sniffing, nothing runs,
916/// and a type a browser would open is an attachment. The site's Worker
917/// sets the same (apps/web/app/lib/content-safety.ts); this keeps the
918/// service safe on its own.
919fn harden(mut response: Response) -> Result<Response> {
920 let headers = response.headers_mut();
921 headers.set("x-content-type-options", "nosniff")?;
922 headers.set("content-security-policy", NOTHING_RUNS)?;
923 if headers.get("content-disposition")?.is_none() && opens_as_document(headers.get("content-type")?.as_deref()) {
924 headers.set("content-disposition", "attachment")?;
925 }
926 Ok(response)
927}
928
929#[event(fetch)]
930async fn fetch(mut request: Request, env: Env, ctx: Context) -> Result<Response> {
931 let packages = Packages::from_env(&env)?;
932 let Some(method) = rpc_method(&request) else {
933 let answer = packages.serve_registry(request, &ctx).await?;
934 return harden(answer);
935 };
936 let body: serde_json::Value = request.json().await?;
937 match method.as_str() {
938 "list_packages" => reply(&packages.list_packages(args(body)?).await?),
939 "get_package" => reply(&packages.get_package(args(body)?).await?),
940 "delete_version" => reply(&packages.delete_version(args(body)?).await?),
941 "delete_package" => reply(&packages.delete_package(args(body)?).await?),
942 "set_package" => reply(&packages.set_package(args(body)?).await?),
943 // A package's settings, deleted packages and versions: settings.rs.
944 "list_versions" => reply(&packages.list_versions(args(body)?).await?),
945 "get_version" => reply(&packages.get_version(args(body)?).await?),
946 "restore_version" => reply(&packages.restore_version(args(body)?).await?),
947 "restore_package" => reply(&packages.restore_package(args(body)?).await?),
948 "deleted_packages" => reply(&packages.deleted_packages(args(body)?).await?),
949 "package_settings" => reply(&packages.package_settings(args(body)?).await?),
950 "set_package_access" => reply(&packages.set_package_access(args(body)?).await?),
951 "remove_package_access" => reply(&packages.remove_package_access(args(body)?).await?),
952 "set_actions_access" => reply(&packages.set_actions_access(args(body)?).await?),
953 "remove_actions_access" => reply(&packages.remove_actions_access(args(body)?).await?),
954 // For billing: what a workspace's packages hold.
955 "storage_all" => reply(&packages.db.storage_all().await?),
956 // Read a repository's Composer package again now, as a push would.
957 "sync_composer" => {
958 let a: SyncComposerArgs = args(body)?;
959 reply(&packages.sync_composer(&a.repo_id).await?)
960 }
961 "storage" => {
962 let a: StorageArgs = args(body)?;
963 let (public_bytes, private_bytes) = packages.db.storage(&a.workspace.to_lowercase()).await?;
964 reply(&PackageStorage { public_bytes, private_bytes })
965 }
966 _ => Response::error("Unknown method", 404),
967 }
968}
969
970/// Every hour: expired uploads are let go, packages and versions deleted
971/// more than 30 days ago are purged, and blobs no version has used for a
972/// day are deleted from the store.
973#[event(scheduled)]
974async fn scheduled(_event: ScheduledEvent, env: Env, _ctx: ScheduleContext) {
975 let packages = match Packages::from_env(&env) {
976 Ok(packages) => packages,
977 Err(error) => {
978 worker::console_error!("packages: the sweep could not start: {error}");
979 return;
980 }
981 };
982 // Before the sweep, so the files of what is purged go with it a day on.
983 match packages.purge(now_ms()).await {
984 Ok((0, 0)) => {}
985 Ok((purged, versions)) => worker::console_log!("packages: purged {purged} deleted packages and {versions} deleted versions"),
986 Err(error) => worker::console_error!("packages: the purge failed: {error}"),
987 }
988 match packages.sweep().await {
989 Ok((0, 0)) => {}
990 Ok((uploads, blobs)) => worker::console_log!("packages: let go of {uploads} uploads and {blobs} blobs"),
991 Err(error) => worker::console_error!("packages: the sweep failed: {error}"),
992 }
993 // Repositories that had a composer.json before the registry did.
994 match packages.composer_backfill().await {
995 Ok(0) => {}
996 Ok(found) => worker::console_log!("packages: found {found} Composer packages"),
997 Err(error) => worker::console_error!("packages: the Composer backfill failed: {error}"),
998 }
999}
1000
1001#[event(queue)]
1002async fn queue(batch: MessageBatch<Event>, env: Env, _ctx: Context) -> Result<()> {
1003 let packages = Packages::from_env(&env)?;
1004 // Each event is acknowledged or retried on its own, so one that fails
1005 // does not run the rest of its batch again.
1006 for message in batch.messages()? {
1007 match packages.on_event(&env, message.body()).await {
1008 Ok(_) => message.ack(),
1009 Err(error) => {
1010 worker::console_error!("packages: event {} failed: {error}", message.body().id);
1011 message.retry();
1012 }
1013 }
1014 }
1015 Ok(())
1016}
1017
1018#[cfg(test)]
1019mod tests {
1020 use super::*;
1021 use serde_json::json;
1022
1023 #[test]
1024 fn publisher_documents_are_downloads() {
1025 for kind in ["application/xml", "text/xml; charset=utf-8", "application/xhtml+xml", "text/html", "image/svg+xml", "application/vnd.foo+xml", "", "text/javascript"] {
1026 assert!(opens_as_document(Some(kind)), "{kind}");
1027 }
1028 assert!(opens_as_document(None));
1029 for kind in [
1030 "application/json",
1031 "text/plain; charset=utf-8",
1032 "application/vnd.oci.image.manifest.v1+json",
1033 "application/vnd.npm.install-v1+json",
1034 "application/octet-stream",
1035 "application/java-archive",
1036 "application/gzip",
1037 "application/pgp-signature",
1038 ] {
1039 assert!(!opens_as_document(Some(kind)), "{kind}");
1040 }
1041 }
1042
1043 fn event(kind: &str, data: serde_json::Value) -> Event {
1044 Event {
1045 id: "evt_1".into(),
1046 kind: kind.into(),
1047 source: "identity".into(),
1048 time: "2026-10-06T00:00:00.000Z".into(),
1049 repo_id: None,
1050 actor: None,
1051 data,
1052 }
1053 }
1054
1055 #[test]
1056 fn downloads_are_counted_by_version_and_written_in_batches() {
1057 let mut counts = (HashMap::new(), 0);
1058 let version = |v: &str| ("pkg_1".to_owned(), Some(v.to_owned()));
1059 // The first is written at once (nothing was written before).
1060 let first = tally(&mut counts, version("ver_1"), 100_000).unwrap();
1061 assert_eq!(first, vec![(version("ver_1"), 1)]);
1062 // Then they are kept until the flush is due, each version apart.
1063 assert!(tally(&mut counts, version("ver_1"), 101_000).is_none());
1064 assert!(tally(&mut counts, version("ver_1"), 102_000).is_none());
1065 assert!(tally(&mut counts, version("ver_2"), 103_000).is_none());
1066 let mut due = tally(&mut counts, ("pkg_1".to_owned(), None), 100_000 + DOWNLOADS_FLUSH_MS).unwrap();
1067 due.sort();
1068 assert_eq!(due, vec![(("pkg_1".to_owned(), None), 1), (version("ver_1"), 2), (version("ver_2"), 1)]);
1069 assert!(counts.0.is_empty());
1070 }
1071
1072 fn package(deleted_at: Option<&str>, workspace_deleted_at: Option<&str>) -> PackageRow {
1073 PackageRow {
1074 id: "pkg_1".into(),
1075 workspace: "acme".into(),
1076 ecosystem: "npm".into(),
1077 name: "web".into(),
1078 repo_id: None,
1079 repo_name: None,
1080 visibility: "private".into(),
1081 description: None,
1082 created_by: "usr_1".into(),
1083 created_at: "2026-10-01T00:00:00.000Z".into(),
1084 updated_at: "2026-10-01T00:00:00.000Z".into(),
1085 downloads: 0,
1086 workspace_deleted_at: workspace_deleted_at.map(str::to_owned),
1087 inherit_access: 1,
1088 deleted_at: deleted_at.map(str::to_owned),
1089 deleted_by: deleted_at.map(|_| "ana".to_owned()),
1090 }
1091 }
1092
1093 #[test]
1094 fn a_deleted_package_is_hidden_and_keeps_its_name_until_the_purge() {
1095 let deleted = package(Some("2026-10-02T09:00:00.000Z"), None);
1096 assert!(deleted.hidden());
1097 let refusal = Packages::hidden_refusal(&deleted);
1098 assert!(refusal.contains("was deleted"), "{refusal}");
1099 assert!(refusal.contains("2026-11-01"), "kept 30 days: {refusal}");
1100 let workspace_gone = package(None, Some("2026-10-02T09:00:00.000Z"));
1101 assert!(workspace_gone.hidden());
1102 assert!(Packages::hidden_refusal(&workspace_gone).contains("workspace acme is deleted"));
1103 assert!(!package(None, None).hidden());
1104 }
1105
1106 #[test]
1107 fn the_purge_takes_what_was_deleted_more_than_thirty_days_ago() {
1108 let now = g1t_contracts::time::parse_rfc3339("2026-11-01T12:00:00.000Z").unwrap();
1109 let before = settings::purge_cutoff(now);
1110 assert_eq!(before, "2026-10-02T12:00:00.000Z");
1111 // Rows are purged when deleted_at < before: a day older goes, a
1112 // second newer stays restorable.
1113 assert!("2026-10-01T12:00:00.000Z" < before.as_str());
1114 assert!("2026-10-02T12:00:01.000Z" > before.as_str());
1115 assert!(settings::restorable("2026-10-02T12:00:01.000Z", now));
1116 assert!(!settings::restorable("2026-10-01T12:00:00.000Z", now));
1117 }
1118
1119 #[test]
1120 fn a_deleted_workspace_is_hidden_and_a_restored_one_shown() {
1121 let protected = g1t_contracts::identity::protected_names(Some("wsp_keep"));
1122 let now = "2026-10-06T12:00:00.000Z";
1123 let deleting = event("workspace.deleting", json!({ "workspaceId": "wsp_1", "slug": "Acme", "by": "ana", "purgeAfter": "x" }));
1124 assert_eq!(
1125 workspace_mark(&deleting, &protected, now),
1126 Some(WorkspaceMark::Hide { workspace: "acme".into(), at: now.into() })
1127 );
1128 let restored = event("workspace.restored", json!({ "workspaceId": "wsp_1", "slug": "acme" }));
1129 assert_eq!(workspace_mark(&restored, &protected, now), Some(WorkspaceMark::Show { workspace: "acme".into() }));
1130 }
1131
1132 #[test]
1133 fn protected_workspaces_and_other_events_are_left_alone() {
1134 let protected = g1t_contracts::identity::protected_names(Some("wsp_keep"));
1135 let now = "2026-10-06T12:00:00.000Z";
1136 let flagon = event("workspace.deleting", json!({ "workspaceId": "wsp_f", "slug": "flagon-io" }));
1137 assert_eq!(workspace_mark(&flagon, &protected, now), None, "always protected");
1138 let by_id = event("workspace.deleting", json!({ "workspaceId": "wsp_keep", "slug": "kept" }));
1139 assert_eq!(workspace_mark(&by_id, &protected, now), None, "protected by id");
1140 let purged = event("workspace.deleted", json!({ "workspaceId": "wsp_1", "slug": "acme" }));
1141 assert_eq!(workspace_mark(&purged, &protected, now), None, "the purge is handled apart");
1142 let nameless = event("workspace.deleting", json!({ "workspaceId": "wsp_1", "slug": "" }));
1143 assert_eq!(workspace_mark(&nameless, &protected, now), None);
1144 }
1145}