g1t/services/packages/src/lib.rs

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