g1t/services/packages/src/lib.rs

712 lines30,601 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 db;
12mod digest;
13mod limits;
14mod manifest;
15mod names;
16mod oci;
17mod quota;
18mod range;
19mod sigv4;
20mod store;
21mod token;
22mod upload;
23
24use std::cell::RefCell;
25use std::collections::HashMap;
26
27use g1t_contracts::audit::{AuditActor, AuditTarget, NewAuditEntry, RecordAuditArgs, Surface};
28use g1t_contracts::credentials::Decision;
29use g1t_contracts::events::{Event, NewEvent, PackageEvent, Publish};
30use g1t_contracts::identity::GitCredentialsArgs;
31use g1t_contracts::packages::*;
32use g1t_contracts::repos::{GetArgs, Repo, RepoPath};
33use g1t_contracts::{FailureCode, Outcome, User, Viewer, new_id};
34use g1t_kit::{args, now_ms, reply, rpc_method};
35use worker::{Context, Env, Fetcher, MessageBatch, Request, Response, Result, ScheduleContext, ScheduledEvent, event};
36
37use access::{Action, LinkedTo, Target};
38use db::{Db, PackageRow};
39use store::{BlobStore, Store};
40
41pub(crate) const SOURCE: &str = "packages";
42/// How large a request body may be: Cloudflare's limit on the zone's plan.
43const DEFAULT_MAX_REQUEST_BYTES: u64 = 100_000_000;
44/// How many packages a listing shows.
45const LIST_LIMIT: u32 = 200;
46const VERSIONS_SHOWN: u32 = 200;
47/// How much one sweep lets go of.
48const SWEEP_BATCH: u32 = 200;
49
50thread_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}
57const DOWNLOADS_FLUSH_MS: u64 = 10_000;
58
59thread_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.
64const ALLOWANCE_TTL_MS: u64 = 5 * 60 * 1000;
65
66/// Who made a registry request, as its audit entries and versions name them.
67pub struct Caller {
68 pub actor: Option<AuditActor>,
69}
70
71/// What access decisions need about a package, owned.
72pub struct TargetOf {
73 workspace: String,
74 repo: Option<(String, String, bool)>,
75 public: bool,
76}
77
78impl 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
99pub 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)]
118pub(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.
128pub(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
144fn not_found<T>() -> Outcome<T> {
145 Outcome::fail(FailureCode::NotFound, "Package not found.")
146}
147
148impl 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)]
619async 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)]
645async 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)]
661async 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)]
670mod 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}