Pick any line to see why it is the way it is: the commit, the pull request and issue it came from, and what the agent was thinking.
| Composer from the workspace's own repositories, and go get from g1t.sh | 1 | //! The Composer registry: `g1t.sh/-/composer/<workspace>/`, one per |
| 2 | //! workspace, built from its repositories (composer.rs says how). | |
| 3 | //! | |
| 4 | //! ```sh | |
| 5 | //! composer config repositories.acme composer https://g1t.sh/-/composer/acme/ | |
| 6 | //! composer config --global --auth http-basic.g1t.sh <you> <g1t token> | |
| 7 | //! composer require acme/lib | |
| 8 | //! ``` | |
| 9 | //! | |
| 10 | //! Nothing is uploaded. A repository's versions are read again when it is | |
| 11 | //! pushed to (`git.push`), restored, renamed or moved, and once for every | |
| 12 | //! repository by the backfill; a request only reads what that left. A zip | |
| 13 | //! of a commit is made the first time it is asked for, and kept by its | |
| 14 | //! digest like every file here. | |
| 15 | ||
| 16 | use std::collections::{HashMap, HashSet}; | |
| 17 | ||
| 18 | use base64::Engine; | |
| 19 | use base64::engine::general_purpose::STANDARD; | |
| 20 | use g1t_contracts::User; | |
| 21 | use g1t_contracts::events::PackageEvent; | |
| 22 | use g1t_contracts::new_id; | |
| 23 | use g1t_contracts::repos::{ | |
| 24 | AllIdsArgs, FileList, IdPage, ListFilesArgs, MAX_LISTED_FILES, MAX_READ_BLOBS, RawBlob, RawBlobsArgs, RawFile, RawFileArgs, RefsArgs, | |
| 25 | RepoRefs, | |
| 26 | }; | |
| 27 | use g1t_kit::now_ms; | |
| 28 | use serde_json::{Value, json}; | |
| 29 | use worker::{Context, Headers, Method, Request, Response, Result, Url}; | |
| 30 | ||
| 31 | use crate::access::{self, Action}; | |
| 32 | use crate::composer::{self, Origin, Version}; | |
| 33 | use crate::db::{NewVersion, PackageRow}; | |
| 34 | use crate::digest::Digest; | |
| 35 | use crate::oci::{Credentials, origin}; | |
| 36 | use crate::store::BlobStore; | |
| 37 | use crate::{Caller, Packages, TargetOf}; | |
| 38 | ||
| 39 | const COMPOSER: &str = "composer"; | |
| 40 | /// The biggest `composer.json` or README read. | |
| 41 | const MAX_FILE_BYTES: u32 = 1024 * 1024; | |
| 42 | /// The most branches and tags a repository's package lists. | |
| 43 | const MAX_BRANCHES: usize = 50; | |
| 44 | const MAX_TAGS: usize = 300; | |
| 45 | /// The most a zip may hold before it is made, in all and per file. | |
| 46 | const MAX_ARCHIVE_BYTES: u64 = 64 * 1024 * 1024; | |
| 47 | const MAX_ARCHIVED_FILE: u32 = 32 * 1024 * 1024; | |
| 48 | /// Repositories the backfill reads each hour. | |
| 49 | const BACKFILL_PAGE: u32 = 25; | |
| 50 | const README_NAMES: [&str; 4] = ["README.md", "readme.md", "README.markdown", "README"]; | |
| 51 | ||
| 52 | fn error(status: u16, message: impl Into<String>) -> Result<Response> { | |
| 53 | let mut response = Response::from_json(&json!({ "status": "error", "message": message.into() }))?.with_status(status); | |
| 54 | if status == 401 { | |
| 55 | response.headers_mut().set("www-authenticate", "Basic realm=\"g1t\"")?; | |
| 56 | } | |
| 57 | Ok(response) | |
| 58 | } | |
| 59 | ||
| 60 | /// One of the registry's endpoints, under `/-/composer/<workspace>/`. | |
| 61 | #[derive(Clone, Debug, PartialEq, Eq)] | |
| 62 | pub enum ComposerRoute { | |
| 63 | Root { workspace: String }, | |
| 64 | /// `p2/<vendor>/<name>.json`, or `~dev.json` for the branches. | |
| 65 | Metadata { workspace: String, name: String, dev: bool }, | |
| 66 | Dist { workspace: String, name: String, commit: String }, | |
| 67 | Downloads { workspace: String }, | |
| 68 | } | |
| 69 | ||
| 70 | pub fn route(path: &str) -> Option<ComposerRoute> { | |
| 71 | let rest = path.strip_prefix("/-/composer/")?; | |
| 72 | let (workspace, rest) = rest.split_once('/')?; | |
| 73 | let workspace = workspace.to_ascii_lowercase(); | |
| 74 | if rest == "packages.json" || rest.is_empty() { | |
| 75 | return Some(ComposerRoute::Root { workspace }); | |
| 76 | } | |
| 77 | if rest == "downloads" { | |
| 78 | return Some(ComposerRoute::Downloads { workspace }); | |
| 79 | } | |
| 80 | if let Some(file) = rest.strip_prefix("p2/") { | |
| 81 | let (name, dev) = match file.strip_suffix("~dev.json") { | |
| 82 | Some(name) => (name, true), | |
| 83 | None => (file.strip_suffix(".json")?, false), | |
| 84 | }; | |
| 85 | return composer::valid_name(name).then(|| ComposerRoute::Metadata { workspace, name: name.to_owned(), dev }); | |
| 86 | } | |
| 87 | let file = rest.strip_prefix("dist/")?; | |
| 88 | let (name, zip) = file.rsplit_once('/')?; | |
| 89 | let commit = zip.strip_suffix(".zip")?; | |
| 90 | (composer::valid_name(name) && commit.len() == 40 && commit.bytes().all(|b| b.is_ascii_hexdigit())).then(|| ComposerRoute::Dist { | |
| 91 | workspace, | |
| 92 | name: name.to_owned(), | |
| 93 | commit: commit.to_ascii_lowercase(), | |
| 94 | }) | |
| 95 | } | |
| 96 | ||
| 97 | /// What a version keeps of its ref, to make its entry from on each read. | |
| 98 | fn stored_metadata(composer_json: &Value, version: &Version, git_ref: &str, default_branch: bool) -> Value { | |
| 99 | json!({ | |
| 100 | "composer": composer_json, | |
| 101 | "version_normalized": version.normalized, | |
| 102 | "ref": git_ref, | |
| 103 | "default_branch": default_branch, | |
| 104 | }) | |
| 105 | } | |
| 106 | ||
| 107 | /// The versions a repository's refs make: `(version, commit, ref, default)`. | |
| 108 | fn wanted_versions(refs: &[g1t_contracts::repos::GitRefEntry], default_branch: &str) -> Vec<(Version, String, String, bool)> { | |
| 109 | let mut branches = Vec::new(); | |
| 110 | let mut tags = Vec::new(); | |
| 111 | for entry in refs { | |
| 112 | if let Some(branch) = entry.name.strip_prefix("refs/heads/") { | |
| 113 | branches.push((composer::branch_version(branch), entry.commit.clone(), entry.name.clone(), branch == default_branch)); | |
| 114 | } else if let Some(tag) = entry.name.strip_prefix("refs/tags/") | |
| 115 | && let Some(version) = composer::tag_version(tag) | |
| 116 | { | |
| 117 | tags.push((version, entry.commit.clone(), entry.name.clone(), false)); | |
| 118 | } | |
| 119 | } | |
| 120 | // The default branch first, then the newest tags. | |
| 121 | branches.sort_by_key(|(_, _, _, default)| !*default); | |
| 122 | branches.truncate(MAX_BRANCHES); | |
| 123 | tags.sort_by(|a, b| b.0.normalized.cmp(&a.0.normalized)); | |
| 124 | tags.truncate(MAX_TAGS); | |
| 125 | let mut seen = HashSet::new(); | |
| 126 | branches.into_iter().chain(tags).filter(|(v, ..)| seen.insert(v.version.clone())).collect() | |
| 127 | } | |
| 128 | ||
| 129 | impl Packages { | |
| 130 | pub async fn composer(&self, request: Request, ctx: &Context) -> Result<Response> { | |
| 131 | let url = request.url()?; | |
| 132 | let Some(route) = route(url.path()) else { | |
| 133 | return error(404, "There is nothing at this address."); | |
| 134 | }; | |
| 135 | match self.composer_route(request, &url, route, ctx).await { | |
| 136 | Ok(response) => Ok(response), | |
| 137 | Err(problem) => { | |
| 138 | worker::console_error!("packages: composer {}: {problem}", url.path()); | |
| 139 | error(500, "Something went wrong on our side. Try again in a moment.") | |
| 140 | } | |
| 141 | } | |
| 142 | } | |
| 143 | ||
| 144 | async fn composer_route(&self, mut request: Request, url: &Url, route: ComposerRoute, ctx: &Context) -> Result<Response> { | |
| 145 | let credentials = self.credentials(&request).await?; | |
| 146 | if matches!(request.method(), Method::Get | Method::Head) | |
| 147 | && let Some(refused) = self.limited(&request, &credentials, "http-basic credentials").await? | |
| 148 | { | |
| 149 | return Ok(refused); | |
| 150 | } | |
| 151 | let viewer = match credentials { | |
| 152 | Credentials::Viewer(viewer) => viewer, | |
| 153 | Credentials::None => None, | |
| 154 | Credentials::Token(_) | Credentials::Bad => { | |
| 155 | return error(401, "The username or token is not right. Use a g1t access token: composer config --auth http-basic.g1t.sh <you> <token>"); | |
| 156 | } | |
| 157 | }; | |
| 158 | match route { | |
| 159 | ComposerRoute::Root { workspace } => self.composer_root(&workspace, viewer.as_ref()).await, | |
| 160 | ComposerRoute::Metadata { workspace, name, dev } => self.composer_metadata(url, &workspace, &name, dev, viewer.as_ref()).await, | |
| 161 | ComposerRoute::Dist { workspace, name, commit } => self.composer_dist(&workspace, &name, &commit, viewer.as_ref(), ctx).await, | |
| 162 | ComposerRoute::Downloads { workspace } => { | |
| 163 | let body: Value = request.json().await.unwrap_or_default(); | |
| 164 | let names: HashSet<&str> = body["downloads"] | |
| 165 | .as_array() | |
| 166 | .map(|list| list.iter().filter_map(|d| d["name"].as_str()).collect()) | |
| 167 | .unwrap_or_default(); | |
| 168 | for name in names.into_iter().take(50) { | |
| 169 | if let Some(package) = self.composer_package(&workspace, name).await? { | |
| 170 | self.count_download(&package.id, ctx); | |
| 171 | } | |
| 172 | } | |
| 173 | Ok(Response::empty()?.with_status(204)) | |
| 174 | } | |
| 175 | } | |
| 176 | } | |
| 177 | ||
| 178 | async fn composer_package(&self, workspace: &str, name: &str) -> Result<Option<PackageRow>> { | |
| 179 | Ok(self.db.package(workspace, COMPOSER, name).await?.filter(|p| !p.hidden())) | |
| 180 | } | |
| 181 | ||
| 182 | /// The answer when `viewer` may not pull `package`, if they may not. | |
| 183 | fn composer_check(&self, viewer: Option<&User>, package: &PackageRow) -> Option<Result<Response>> { | |
| 184 | let target = TargetOf::package(package); | |
| 185 | if access::decide(viewer, &target.view(), Action::Pull).allowed { | |
| 186 | return None; | |
| 187 | } | |
| 188 | Some(if viewer.is_none() { | |
| 189 | error(401, "Sign in to install this package: composer config --auth http-basic.g1t.sh <you> <g1t token>") | |
| 190 | } else { | |
| 191 | error(404, "Not found: no such package, or you cannot see it.") | |
| 192 | }) | |
| 193 | } | |
| 194 | ||
| 195 | async fn composer_root(&self, workspace: &str, viewer: Option<&User>) -> Result<Response> { | |
| 196 | let rows = self.db.list(workspace, Some(COMPOSER), None, None, 1000).await?; | |
| 197 | let available: Vec<String> = rows | |
| 198 | .iter() | |
| 199 | .filter(|row| access::decide(viewer, &TargetOf::package(&row.package).view(), Action::Pull).allowed) | |
| 200 | .map(|row| row.package.name.clone()) | |
| 201 | .collect(); | |
| 202 | Response::from_json(&composer::root(workspace, &available)) | |
| 203 | } | |
| 204 | ||
| 205 | async fn composer_metadata(&self, url: &Url, workspace: &str, name: &str, dev: bool, viewer: Option<&User>) -> Result<Response> { | |
| 206 | let Some(package) = self.composer_package(workspace, name).await? else { | |
| 207 | return error(404, format!("There is no package {name} in {workspace}.")); | |
| 208 | }; | |
| 209 | if let Some(refusal) = self.composer_check(viewer, &package) { | |
| 210 | return refusal; | |
| 211 | } | |
| 212 | let base = origin(url); | |
| 213 | let repo = package.repo_name.clone().unwrap_or_default(); | |
| 214 | let git_url = format!("{base}/{}/{repo}.git", package.workspace); | |
| 215 | let mut entries = Vec::new(); | |
| 216 | for row in self.db.versions(&package.id, 1000).await? { | |
| 217 | let meta = row.meta(); | |
| 218 | let version = Version { | |
| 219 | version: row.version.clone(), | |
| 220 | normalized: meta["version_normalized"].as_str().unwrap_or(&row.version).to_owned(), | |
| 221 | }; | |
| 222 | if version.is_dev() != dev { | |
| 223 | continue; | |
| 224 | } | |
| 225 | let dist_url = format!("{base}/-/composer/{}/dist/{name}/{}.zip", package.workspace, row.digest); | |
| 226 | let origin = Origin { | |
| 227 | git_url: &git_url, | |
| 228 | dist_url: &dist_url, | |
| 229 | commit: &row.digest, | |
| 230 | default_branch: meta["default_branch"].as_bool().unwrap_or(false), | |
| 231 | }; | |
| 232 | entries.push(composer::version_entry(&meta["composer"], name, &version, &origin)); | |
| 233 | } | |
| 234 | let mut response = Response::from_json(&composer::p2(name, &entries))?; | |
| 235 | response.headers_mut().set("last-modified", &package.updated_at)?; | |
| 236 | Ok(response) | |
| 237 | } | |
| 238 | ||
| 239 | async fn composer_dist(&self, workspace: &str, name: &str, commit: &str, viewer: Option<&User>, ctx: &Context) -> Result<Response> { | |
| 240 | let Some(package) = self.composer_package(workspace, name).await? else { | |
| 241 | return error(404, format!("There is no package {name} in {workspace}.")); | |
| 242 | }; | |
| 243 | if let Some(refusal) = self.composer_check(viewer, &package) { | |
| 244 | return refusal; | |
| 245 | } | |
| 246 | // Only the commits of its versions: a zip is never made of any | |
| 247 | // other commit of the repository. | |
| 248 | if self.db.version_by_digest(&package.id, commit).await?.is_none() { | |
| 249 | return error(404, format!("{commit} is not a version of {name}.")); | |
| 250 | } | |
| 251 | let blob = match self.db.dist_for_commit(&package.id, commit).await? { | |
| 252 | Some(blob) => blob, | |
| 253 | None => match self.build_dist(&package, commit).await? { | |
| 254 | Ok(blob) => blob, | |
| 255 | Err(refusal) => return refusal, | |
| 256 | }, | |
| 257 | }; | |
| 258 | let Some(got) = self.store.get(&blob.object_key, None).await? else { | |
| 259 | return error(404, "The archive is missing. Try again."); | |
| 260 | }; | |
| 261 | self.count_download(&package.id, ctx); | |
| 262 | let headers = Headers::new(); | |
| 263 | headers.set("content-type", "application/zip")?; | |
| 264 | headers.set("content-length", &blob.size.to_string())?; | |
| 265 | headers.set("cache-control", "max-age=31536000")?; | |
| 266 | Ok(Response::from_body(got.body)?.with_headers(headers)) | |
| 267 | } | |
| 268 | ||
| 269 | /// Makes the zip of a commit: its files but those `.gitattributes` | |
| 270 | /// marks `export-ignore`, as `git archive` would leave them out. | |
| 271 | async fn build_dist(&self, package: &PackageRow, commit: &str) -> Result<std::result::Result<crate::db::BlobRow, Result<Response>>> { | |
| 272 | let Some(repo_id) = package.repo_id.clone() else { | |
| 273 | return Ok(Err(error(404, "This package has no repository."))); | |
| 274 | }; | |
| 275 | let listed: FileList = g1t_kit::call( | |
| 276 | &self.repos, | |
| 277 | "list_files", | |
| 278 | &ListFilesArgs { repo_id: repo_id.clone(), git_ref: Some(commit.to_owned()), skip_dirs: Vec::new(), limit: MAX_LISTED_FILES }, | |
| 279 | ) | |
| 280 | .await?; | |
| 281 | if listed.truncated { | |
| 282 | return Ok(Err(error(507, format!("The commit has more than {MAX_LISTED_FILES} files, too many for an archive. Install from source: composer install --prefer-source")))); | |
| 283 | } | |
| 284 | let attributes = self.repo_file(&repo_id, commit, ".gitattributes").await?; | |
| 285 | let ignores = attributes.map(|text| composer::export_ignores(&String::from_utf8_lossy(&text))).unwrap_or_default(); | |
| 286 | let files: Vec<(String, String)> = listed | |
| 287 | .files | |
| 288 | .into_iter() | |
| 289 | .filter_map(|file| Some((file.path, file.hash?))) | |
| 290 | .filter(|(path, _)| !composer::ignored(&ignores, path)) | |
| 291 | .collect(); | |
| 292 | let mut bytes_of: HashMap<String, Vec<u8>> = HashMap::new(); | |
| 293 | let unique: Vec<String> = files.iter().map(|(_, hash)| hash.clone()).collect::<HashSet<_>>().into_iter().collect(); | |
| 294 | let mut total = 0u64; | |
| 295 | for chunk in unique.chunks(MAX_READ_BLOBS) { | |
| 296 | let read: Vec<RawBlob> = g1t_kit::call( | |
| 297 | &self.repos, | |
| 298 | "raw_blobs", | |
| 299 | &RawBlobsArgs { repo_id: repo_id.clone(), hashes: chunk.to_vec(), max_bytes: MAX_ARCHIVED_FILE }, | |
| 300 | ) | |
| 301 | .await?; | |
| 302 | for blob in read { | |
| 303 | total += blob.size; | |
| 304 | if total > MAX_ARCHIVE_BYTES || (blob.data.is_none() && blob.size > 0) { | |
| 305 | return Ok(Err(error( | |
| 306 | 507, | |
| 307 | format!("The commit is too large for an archive (over {} MB). Install from source: composer install --prefer-source", MAX_ARCHIVE_BYTES / 1_048_576), | |
| 308 | ))); | |
| 309 | } | |
| 310 | let data = blob.data.as_deref().map(|d| STANDARD.decode(d).unwrap_or_default()).unwrap_or_default(); | |
| 311 | bytes_of.insert(blob.hash, data); | |
| 312 | } | |
| 313 | } | |
| 314 | let entries: Vec<(String, Vec<u8>)> = files | |
| 315 | .into_iter() | |
| 316 | .map(|(path, hash)| { | |
| 317 | let data = bytes_of.get(&hash).cloned().unwrap_or_default(); | |
| 318 | (path, data) | |
| 319 | }) | |
| 320 | .collect(); | |
| 321 | let zip = composer::zip(&entries); | |
| 322 | let digest = Digest::of(&zip); | |
| 323 | let size = zip.len() as u64; | |
| 324 | let now = now_ms(); | |
| 325 | if self.db.blob(&digest).await?.is_none() { | |
| 326 | self.store.put(&digest.object_key(), zip).await?; | |
| 327 | } | |
| 328 | self.db.keep_blob(&package.id, &digest, size, Some("application/zip"), &digest.object_key(), now).await?; | |
| 329 | self.db.add_dist(&package.id, commit, &digest, size).await?; | |
| 330 | self.db.measure(&package.workspace).await?; | |
| 331 | Ok(Ok(crate::db::BlobRow { digest: digest.to_string(), size, media_type: Some("application/zip".into()), object_key: digest.object_key() })) | |
| 332 | } | |
| 333 | ||
| 334 | /// A file of a repository at a commit, if it is there and not large. | |
| 335 | async fn repo_file(&self, repo_id: &str, git_ref: &str, path: &str) -> Result<Option<Vec<u8>>> { | |
| 336 | let file: Option<RawFile> = g1t_kit::call( | |
| 337 | &self.repos, | |
| 338 | "raw_file", | |
| 339 | &RawFileArgs { repo_id: repo_id.to_owned(), git_ref: git_ref.to_owned(), path: path.to_owned(), max_bytes: MAX_FILE_BYTES }, | |
| 340 | ) | |
| 341 | .await?; | |
| 342 | Ok(file.and_then(|file| STANDARD.decode(file.data).ok())) | |
| 343 | } | |
| 344 | ||
| 345 | /// Reads a repository's Composer package again from its refs: makes it | |
| 346 | /// when its default branch gained a `composer.json`, records new and | |
| 347 | /// moved versions, lets go of deleted ones, and deletes the package | |
| 348 | /// when the repository stopped being one. Says whether it is one. | |
| 349 | pub(crate) async fn sync_composer(&self, repo_id: &str) -> Result<bool> { | |
| 350 | let found: Option<RepoRefs> = g1t_kit::call(&self.repos, "refs", &RefsArgs { repo_id: repo_id.to_owned() }).await?; | |
| 351 | let existing = self.db.package_for_repo(repo_id, COMPOSER).await?; | |
| 352 | let Some(RepoRefs { repo, refs }) = found else { | |
| 353 | if let Some(package) = existing { | |
| 354 | self.drop_composer(&package).await?; | |
| 355 | } | |
| 356 | return Ok(false); | |
| 357 | }; | |
| 358 | let workspace = repo.namespace.to_lowercase(); | |
| 359 | if self.db.workspace_hidden(&workspace).await? { | |
| 360 | return Ok(false); | |
| 361 | } | |
| 362 | let default = refs.iter().find(|r| r.name == format!("refs/heads/{}", repo.default_branch)).map(|r| r.commit.clone()); | |
| 363 | let manifest = match &default { | |
| 364 | Some(commit) => self.composer_json(repo_id, commit).await?, | |
| 365 | None => None, | |
| 366 | }; | |
| 367 | let Some((name, root_manifest)) = manifest else { | |
| 368 | if let Some(package) = existing { | |
| 369 | self.drop_composer(&package).await?; | |
| 370 | } | |
| 371 | return Ok(false); | |
| 372 | }; | |
| 373 | // A repository moved to another workspace takes its package along. | |
| 374 | let existing = match existing { | |
| 375 | Some(package) if package.workspace != workspace => { | |
| 376 | self.drop_composer(&package).await?; | |
| 377 | None | |
| 378 | } | |
| 379 | other => other, | |
| 380 | }; | |
| 381 | let now = now_ms(); | |
| 382 | let package = match existing { | |
| 383 | Some(package) if package.name == name => package, | |
| 384 | Some(package) => { | |
| 385 | if self.db.package(&workspace, COMPOSER, &name).await?.is_some() { | |
| 386 | worker::console_error!("packages: {workspace}/{} names {name}, which another repository has", repo.name); | |
| 387 | package | |
| 388 | } else { | |
| 389 | self.db.rename_package(&package.id, &name, now).await?; | |
| 390 | PackageRow { name: name.clone(), ..package } | |
| 391 | } | |
| 392 | } | |
| 393 | None => { | |
| 394 | if let Some(other) = self.db.package(&workspace, COMPOSER, &name).await? | |
| 395 | && other.repo_id.as_deref() != Some(repo_id) | |
| 396 | { | |
| 397 | worker::console_error!("packages: {workspace}/{} names {name}, which another repository has", repo.name); | |
| 398 | return Ok(false); | |
| 399 | } | |
| 400 | self.db | |
| 401 | .create_package(&new_id("pkg", now), &workspace, COMPOSER, &name, Some((&repo.id, &repo.name, repo.is_private)), "g1t", now) | |
| 402 | .await? | |
| 403 | } | |
| 404 | }; | |
| 405 | ||
| 406 | let caller = Caller { actor: None }; | |
| 407 | let wanted = wanted_versions(&refs, &repo.default_branch); | |
| 408 | let stored = self.db.versions(&package.id, 1000).await?; | |
| 409 | let mut manifests: HashMap<String, Option<Value>> = HashMap::new(); | |
| 410 | if let Some(commit) = &default { | |
| 411 | manifests.insert(commit.clone(), Some(root_manifest.clone())); | |
| 412 | } | |
| 413 | let mut changed = false; | |
| 414 | for (version, commit, git_ref, is_default) in &wanted { | |
| 415 | let current = stored.iter().find(|row| row.version == version.version); | |
| 416 | if let Some(row) = current | |
| 417 | && row.digest == *commit | |
| 418 | && row.meta()["default_branch"].as_bool().unwrap_or(false) == *is_default | |
| 419 | { | |
| 420 | continue; | |
| 421 | } | |
| 422 | if !manifests.contains_key(commit) { | |
| 423 | let read = self.composer_json(repo_id, commit).await?.map(|(_, json)| json); | |
| 424 | manifests.insert(commit.clone(), read); | |
| 425 | } | |
| 426 | // A ref without a composer.json of its own is not a version. | |
| 427 | let Some(Some(json)) = manifests.get(commit) else { continue }; | |
| 428 | self.db | |
| 429 | .replace_version( | |
| 430 | &NewVersion { | |
| 431 | id: new_id("ver", now), | |
| 432 | package_id: package.id.clone(), | |
| 433 | version: version.version.clone(), | |
| 434 | digest: commit.clone(), | |
| 435 | size: 0, | |
| 436 | metadata: stored_metadata(json, version, git_ref, *is_default).to_string(), | |
| 437 | subject: None, | |
| 438 | published_by: None, | |
| 439 | files: Vec::new(), | |
| 440 | }, | |
| 441 | now, | |
| 442 | ) | |
| 443 | .await?; | |
| 444 | changed = true; | |
| 445 | if current.is_none() { | |
| 446 | let event = PackageEvent { version: Some(version.version.clone()), digest: Some(commit.clone()), ..self.event_of(&package) }; | |
| 447 | self.announce("package.published", &package, event, &caller).await; | |
| 448 | } | |
| 449 | } | |
| 450 | let kept: HashSet<&str> = wanted.iter().map(|(v, ..)| v.version.as_str()).collect(); | |
| 451 | for row in stored.iter().filter(|row| !kept.contains(row.version.as_str())) { | |
| 452 | self.db.delete_version(&row.id).await?; | |
| 453 | let event = PackageEvent { version: Some(row.version.clone()), digest: Some(row.digest.clone()), ..self.event_of(&package) }; | |
| 454 | self.announce("package.version_deleted", &package, event, &caller).await; | |
| 455 | changed = true; | |
| 456 | } | |
| 457 | if let Some(commit) = &default { | |
| 458 | self.composer_readme(&package, repo_id, commit, &root_manifest, now).await?; | |
| 459 | } | |
| 460 | if changed { | |
| 461 | self.db.measure(&workspace).await?; | |
| 462 | } | |
| 463 | Ok(true) | |
| 464 | } | |
| 465 | ||
| 466 | /// The package's README and description, from the default branch. | |
| 467 | async fn composer_readme(&self, package: &PackageRow, repo_id: &str, commit: &str, manifest: &Value, now: u64) -> Result<()> { | |
| 468 | let mut readme = None; | |
| 469 | for name in README_NAMES { | |
| 470 | if let Some(bytes) = self.repo_file(repo_id, commit, name).await? { | |
| 471 | readme = Some(bytes); | |
| 472 | break; | |
| 473 | } | |
| 474 | } | |
| 475 | let digest = match readme.filter(|b| !b.is_empty()) { | |
| 476 | Some(bytes) => { | |
| 477 | let digest = Digest::of(&bytes); | |
| 478 | if self.db.blob(&digest).await?.is_none() { | |
| 479 | self.store.put(&digest.object_key(), bytes.clone()).await?; | |
| 480 | } | |
| 481 | self.db.keep_blob(&package.id, &digest, bytes.len() as u64, Some("text/markdown"), &digest.object_key(), now).await?; | |
| 482 | Some(digest.to_string()) | |
| 483 | } | |
| 484 | None => None, | |
| 485 | }; | |
| 486 | self.db.set_readme(&package.id, digest.as_deref(), manifest["description"].as_str(), now).await | |
| 487 | } | |
| 488 | ||
| 489 | /// A commit's `composer.json`, when it has one naming a valid package. | |
| 490 | async fn composer_json(&self, repo_id: &str, commit: &str) -> Result<Option<(String, Value)>> { | |
| 491 | let Some(bytes) = self.repo_file(repo_id, commit, "composer.json").await? else { | |
| 492 | return Ok(None); | |
| 493 | }; | |
| 494 | let Ok(json) = serde_json::from_slice::<Value>(&bytes) else { | |
| 495 | return Ok(None); | |
| 496 | }; | |
| 497 | let Some(name) = json["name"].as_str().map(str::to_lowercase).filter(|n| composer::valid_name(n)) else { | |
| 498 | return Ok(None); | |
| 499 | }; | |
| 500 | Ok(Some((name, json))) | |
| 501 | } | |
| 502 | ||
| 503 | async fn drop_composer(&self, package: &PackageRow) -> Result<()> { | |
| 504 | self.db.delete_package(&package.id).await?; | |
| 505 | self.db.measure(&package.workspace).await?; | |
| 506 | self.announce("package.deleted", package, self.event_of(package), &Caller { actor: None }).await; | |
| 507 | Ok(()) | |
| 508 | } | |
| 509 | ||
| 510 | /// Deletes the Composer package built from a deleted repository. | |
| 511 | pub(crate) async fn composer_repo_gone(&self, repo_id: &str) -> Result<()> { | |
| 512 | if let Some(package) = self.db.package_for_repo(repo_id, COMPOSER).await? { | |
| 513 | self.drop_composer(&package).await?; | |
| 514 | } | |
| 515 | Ok(()) | |
| 516 | } | |
| 517 | ||
| 518 | /// Reads a page of repositories the backfill has not yet, until it has | |
| 519 | /// read them all once. Says how many were packages. | |
| 520 | pub(crate) async fn composer_backfill(&self) -> Result<u32> { | |
| 521 | let (after, finished) = self.db.backfill().await?; | |
| 522 | if finished { | |
| 523 | return Ok(0); | |
| 524 | } | |
| 525 | let page: IdPage = g1t_kit::call(&self.repos, "all_ids", &AllIdsArgs { after: after.clone(), limit: BACKFILL_PAGE }).await?; | |
| 526 | let mut found = 0; | |
| 527 | for id in &page.ids { | |
| 528 | match self.sync_composer(id).await { | |
| 529 | Ok(true) => found += 1, | |
| 530 | Ok(false) => {} | |
| 531 | Err(error) => worker::console_error!("packages: composer backfill of {id}: {error}"), | |
| 532 | } | |
| 533 | } | |
| 534 | let last = page.ids.last().cloned().or(after); | |
| 535 | self.db.set_backfill(last.as_deref(), page.next.is_none(), now_ms()).await?; | |
| 536 | Ok(found) | |
| 537 | } | |
| 538 | } | |
| 539 | ||
| 540 | #[cfg(test)] | |
| 541 | mod tests { | |
| 542 | use super::*; | |
| 543 | use g1t_contracts::repos::GitRefEntry; | |
| 544 | ||
| 545 | #[test] | |
| 546 | fn every_endpoint_is_routed() { | |
| 547 | let commit = "a".repeat(40); | |
| 548 | assert_eq!(route("/-/composer/acme/packages.json"), Some(ComposerRoute::Root { workspace: "acme".into() })); | |
| 549 | assert_eq!(route("/-/composer/acme/"), Some(ComposerRoute::Root { workspace: "acme".into() })); | |
| 550 | assert_eq!( | |
| 551 | route("/-/composer/acme/p2/acme/lib.json"), | |
| 552 | Some(ComposerRoute::Metadata { workspace: "acme".into(), name: "acme/lib".into(), dev: false }) | |
| 553 | ); | |
| 554 | assert_eq!( | |
| 555 | route("/-/composer/acme/p2/acme/lib~dev.json"), | |
| 556 | Some(ComposerRoute::Metadata { workspace: "acme".into(), name: "acme/lib".into(), dev: true }) | |
| 557 | ); | |
| 558 | assert_eq!( | |
| 559 | route(&format!("/-/composer/acme/dist/acme/lib/{commit}.zip")), | |
| 560 | Some(ComposerRoute::Dist { workspace: "acme".into(), name: "acme/lib".into(), commit: commit.clone() }) | |
| 561 | ); | |
| 562 | assert_eq!(route("/-/composer/acme/downloads"), Some(ComposerRoute::Downloads { workspace: "acme".into() })); | |
| 563 | assert_eq!(route("/-/composer/acme/p2/Acme/lib.json"), None); | |
| 564 | assert_eq!(route("/-/composer/acme/dist/acme/lib/short.zip"), None); | |
| 565 | assert_eq!(route("/-/composer/acme"), None); | |
| 566 | } | |
| 567 | ||
| 568 | #[test] | |
| 569 | fn a_repositorys_refs_make_its_versions() { | |
| 570 | let entry = |name: &str, commit: &str| GitRefEntry { name: name.into(), commit: commit.into() }; | |
| 571 | let refs = [ | |
| 572 | entry("refs/heads/feature", "f"), | |
| 573 | entry("refs/heads/main", "m"), | |
| 574 | entry("refs/tags/v1.0.0", "a"), | |
| 575 | entry("refs/tags/v1.1.0", "b"), | |
| 576 | entry("refs/tags/nightly", "n"), | |
| 577 | ]; | |
| 578 | let wanted = wanted_versions(&refs, "main"); | |
| 579 | let names: Vec<(&str, &str, bool)> = wanted.iter().map(|(v, c, _, d)| (v.version.as_str(), c.as_str(), *d)).collect(); | |
| 580 | assert_eq!( | |
| 581 | names, | |
| 582 | [("dev-main", "m", true), ("dev-feature", "f", false), ("v1.1.0", "b", false), ("v1.0.0", "a", false)], | |
| 583 | "the default branch first, newest tags next, tags that are not versions left out" | |
| 584 | ); | |
| 585 | } | |
| 586 | } |