diff --git a/.github/scripts/start-registry.sh b/.github/scripts/start-registry.sh new file mode 100755 index 0000000..c2686a7 --- /dev/null +++ b/.github/scripts/start-registry.sh @@ -0,0 +1,25 @@ +#!/usr/bin/env bash +# Starts the registry with `wrangler dev` in the background, using the test configuration (user +# "hello", password "world"), and waits until it answers. +# +# start-registry.sh PORT STATE_DIR +set -euo pipefail +port=$1 +state=$2 +mkdir -p "$state" + +WRANGLER_SEND_METRICS=false nohup pnpm exec wrangler dev --config test/wrangler.test.jsonc --env dev \ + --ip 127.0.0.1 --port "$port" --inspector-port 0 --persist-to "$state/r2" >"$state/wrangler.log" 2>&1 & +echo $! >"$state/wrangler.pid" + +for _ in $(seq 1 120); do + if [ "$(curl -s -o /dev/null -w '%{http_code}' "http://127.0.0.1:$port/v2/")" = "401" ]; then + echo "registry is up on port $port" + exit 0 + fi + sleep 1 +done + +echo "registry did not start" >&2 +cat "$state/wrangler.log" >&2 +exit 1 diff --git a/.github/workflows/conformance.yml b/.github/workflows/conformance.yml new file mode 100644 index 0000000..6b08719 --- /dev/null +++ b/.github/workflows/conformance.yml @@ -0,0 +1,145 @@ +name: Conformance + +# Runs the OCI distribution-spec conformance suite against the registry served by `wrangler dev`. + +on: + pull_request: + workflow_dispatch: + +permissions: + contents: read + +jobs: + conformance: + name: OCI conformance (v1.1.1) + runs-on: ubuntu-latest + timeout-minutes: 20 + steps: + - name: Checkout + uses: actions/checkout@v6 + - name: Checkout the conformance suite + uses: actions/checkout@v6 + with: + repository: opencontainers/distribution-spec + ref: v1.1.1 + path: distribution-spec + persist-credentials: false + - name: Install pnpm + uses: pnpm/action-setup@v6 + - name: Use Node + uses: actions/setup-node@v6 + with: + node-version: 24 + cache: "pnpm" + - name: Use Go + uses: actions/setup-go@v6 + with: + go-version-file: distribution-spec/conformance/go.mod + cache-dependency-path: distribution-spec/conformance/go.sum + + - run: pnpm install + - name: Build the conformance suite + working-directory: distribution-spec/conformance + run: go test -c -o "$RUNNER_TEMP/conformance.test" + - name: Start the registry + run: ./.github/scripts/start-registry.sh 5000 "$RUNNER_TEMP/registry" + - name: Run the conformance suite + env: + OCI_ROOT_URL: http://127.0.0.1:5000 + OCI_NAMESPACE: conformance/repo1 + OCI_CROSSMOUNT_NAMESPACE: conformance/repo2 + OCI_USERNAME: hello + OCI_PASSWORD: world + OCI_TEST_PULL: 1 + OCI_TEST_PUSH: 1 + OCI_TEST_CONTENT_DISCOVERY: 1 + OCI_TEST_CONTENT_MANAGEMENT: 1 + OCI_HIDE_SKIPPED_WORKFLOWS: 0 + OCI_DEBUG: 0 + run: | + mkdir -p "$RUNNER_TEMP/results" + cd "$RUNNER_TEMP/results" + OCI_REPORT_DIR="$RUNNER_TEMP/results" "$RUNNER_TEMP/conformance.test" + - name: Show the registry log + if: failure() + run: tail -n 200 "$RUNNER_TEMP/registry/wrangler.log" + - name: Upload the report + if: always() + uses: actions/upload-artifact@v6 + with: + name: conformance-v1.1.1 + path: ${{ runner.temp }}/results + if-no-files-found: ignore + + conformance-main: + # The suite on distribution-spec's main branch is still changing, so this job only reports: a + # failing suite is shown in the job summary but does not fail the job. + name: OCI conformance (main, non-blocking) + runs-on: ubuntu-latest + timeout-minutes: 20 + steps: + - name: Checkout + uses: actions/checkout@v6 + - name: Checkout the conformance suite + uses: actions/checkout@v6 + with: + repository: opencontainers/distribution-spec + ref: main + path: distribution-spec + persist-credentials: false + - name: Install pnpm + uses: pnpm/action-setup@v6 + - name: Use Node + uses: actions/setup-node@v6 + with: + node-version: 24 + cache: "pnpm" + - name: Use Go + uses: actions/setup-go@v6 + with: + go-version-file: distribution-spec/conformance/go.mod + cache-dependency-path: distribution-spec/conformance/go.sum + + - run: pnpm install + - name: Build the conformance suite + working-directory: distribution-spec/conformance + run: go build -o "$RUNNER_TEMP/conformance" . + - name: Start the registry + run: ./.github/scripts/start-registry.sh 5000 "$RUNNER_TEMP/registry" + - name: Run the conformance suite + id: suite + continue-on-error: true + shell: bash + env: + OCI_VERSION: "1.1" + OCI_REGISTRY: 127.0.0.1:5000 + OCI_TLS: disabled + OCI_REPO1: conformance/repo1 + OCI_REPO2: conformance/repo2 + OCI_USERNAME: hello + OCI_PASSWORD: world + # sha512 digests are not supported yet + OCI_DATA_SHA512: "false" + OCI_LOG: error + run: | + mkdir -p "$RUNNER_TEMP/results" + cd "$RUNNER_TEMP/results" + OCI_RESULTS_DIR="$RUNNER_TEMP/results" "$RUNNER_TEMP/conformance" | tee "$RUNNER_TEMP/results/output.txt" + - name: Summarize + if: always() + env: + OUTCOME: ${{ steps.suite.outcome }} + run: | + { + echo "### OCI conformance, distribution-spec main: $OUTCOME" + echo '```' + sed -n '/OCI Conformance Result/,$p' "$RUNNER_TEMP/results/output.txt" 2>/dev/null || true + echo '```' + } >> "$GITHUB_STEP_SUMMARY" + - name: Upload the report + if: always() + uses: actions/upload-artifact@v6 + with: + name: conformance-main + path: ${{ runner.temp }}/results + if-no-files-found: ignore diff --git a/README.md b/README.md index d6759c8..46cd1dc 100644 --- a/README.md +++ b/README.md @@ -99,6 +99,52 @@ docker rmi ubuntu:latest $REGISTRY_URL/ubuntu:latest docker pull $REGISTRY_URL/ubuntu:latest ``` +### Allowing anonymous pulls + +Set `ANONYMOUS_PULL_REPOSITORIES` to a comma or space separated list of repository names that can be pulled without +credentials. `*` matches any characters, including `/`, so `public/*` allows every repository under `public/` and `*` +allows all of them. Requests without an `Authorization` header can then read the manifests, blobs, tags and referrers +of those repositories. Everything else still needs credentials: pushes, deletes, uploads, `/v2/_catalog` and garbage +collection. `/v2/` keeps answering `401` so that clients still log in before they push, requests with wrong +credentials are still refused, and anonymous requests never use the pull fallback below. + +### Protecting immutable release tags + +Set `IMMUTABLE_TAG_PATTERN` under `[env.production.vars]` to a JavaScript regular expression that must match the +entire protected tag. For example, this protects strict `vX.Y.Z` releases while leaving `latest` mutable: + +```toml +IMMUTABLE_TAG_PATTERN = 'v(?:0|[1-9]\d*)\.(?:0|[1-9]\d*)\.(?:0|[1-9]\d*)' +``` + +Protected tags are created with an atomic conditional R2 write. Retrying the same manifest digest is idempotent; +attempting to assign a different digest returns `409` with the OCI `DENIED` error code. Protected tags cannot be +deleted directly. While the policy is enabled, the API rejects every +delete-by-digest request because alias discovery and digest deletion cannot be made atomic across R2 keys. Delete an +unprotected tag by name and let untagged garbage collection remove its content. Direct blob deletion is also disabled +because deleting a referenced layer or config would make a protected release unpullable. An invalid expression fails +manifest writes before any manifest object is stored. + +The policy is enforced at the Worker API boundary. To preserve the invariant, restrict direct R2 write access and +route registry writes through this Worker. + +### Disabling deletion + +Set `DISABLE_DELETE = "true"` to make the registry append-only. Deleting manifests (by tag or by digest), deleting +blobs and garbage collection (`POST /v2//gc`) then answer `405 Method Not Allowed` with the OCI `UNSUPPORTED` +error code. Pushing, and moving a tag that is not protected by `IMMUTABLE_TAG_PATTERN`, keep working. Cancelling an +upload in progress is not affected, because it only removes temporary upload state. + +### Using buckets with retention rules + +Blobs, manifests stored under their digest and referrer entries are written once and never overwritten: a push of +content that already exists leaves the stored object alone. The registry therefore works on an R2 bucket whose +content keys are protected by [bucket locks](https://developers.cloudflare.com/r2/buckets/bucket-locks/) or other +retention rules. Those keys are `/blobs/`, `/manifests/sha256:` and +`/_referrers//`. Tags (`/manifests/`) and upload state +are rewritten and deleted, so keep them outside such rules, and set `DISABLE_DELETE` so that deletes fail cleanly +instead of hitting the lock. + ### Configuring Pull fallback You can configure the R2 registry to fallback to another registry if @@ -163,6 +209,41 @@ REGISTRIES_JSON = "[{ \"registry\": \"https://index.docker.io/\" }]" You can also set your `docker.io` credentials in the configuration to not have any rate-limiting. +### Using the registry from another Worker + +The package can be a dependency of another Worker that does its own routing and hands registry requests to the +registry. Wrangler bundles the TypeScript sources directly, so there is no build step. Pin a commit: + +```jsonc +// package.json of your Worker +"dependencies": { + "r2-registry": "github:cloudflare/serverless-registry#" +} +``` + +```ts +import registry, { type RegistryEnv } from "r2-registry"; + +interface Env extends RegistryEnv { + // your own bindings +} + +export default { + async fetch(request, env, ctx) { + const { pathname } = new URL(request.url); + if (pathname === "/v2" || pathname.startsWith("/v2/")) { + return registry.fetch(request, env, ctx); + } + return new Response("Not Found", { status: 404 }); + }, +} satisfies ExportedHandler; +``` + +`registry.fetch(request, env, ctx)` takes the same bindings and variables as a standalone deployment (`RegistryEnv`): +an R2 bucket bound as `REGISTRY`, and the authentication variables described above. The Worker needs the +`nodejs_compat` compatibility flag. The registry only answers paths under `/v2/`, and it uses the request URL for +authentication challenges and upload locations, so pass the request through with its path unchanged. + ### Known limitations Right now there is some limitations with this container registry. diff --git a/index.ts b/index.ts index 08dacdb..22a5065 100644 --- a/index.ts +++ b/index.ts @@ -1,5 +1,8 @@ /** * The core server that runs on a Cloudflare worker. + * + * Another Worker can import this module and delegate requests to it, see "Using the registry from + * another Worker" in the README. */ import { Router } from "itty-router"; @@ -8,14 +11,19 @@ import v2Router from "./src/router"; import { authenticationMethodFromEnv } from "./src/authentication-method"; import { Registry } from "./src/registry/registry"; import { R2Registry } from "./src/registry/r2"; +import { anonymousPullRepository } from "./src/anonymous"; // A full compatibility mode means that the r2 registry will try its best to // help the client on the layer push. See how we let the client push layers with chunked uploads for more information. -type PushCompatibilityMode = "full" | "none"; +export type PushCompatibilityMode = "full" | "none"; -export interface Env { +/** + * The bindings and variables the registry reads. A Worker that embeds the registry passes an object + * of this shape, usually its own env, to `fetch`. + */ +export interface RegistryEnv { REGISTRY: R2Bucket; - ENVIRONMENT: string; + ENVIRONMENT?: string; JWT_REGISTRY_TOKENS_PUBLIC_KEY?: string; USERNAME?: string; PASSWORD?: string; @@ -23,7 +31,18 @@ export interface Env { READONLY_PASSWORD?: string; PUSH_COMPATIBILITY_MODE?: PushCompatibilityMode; REGISTRIES_JSON?: string; // should be in the format of RegistryConfiguration[]; + // Tags matching this regular expression (the whole tag) are create-only, see src/registry/tag-policy.ts + IMMUTABLE_TAG_PATTERN?: string; + // Set to "true" to refuse every delete: manifests, blobs and garbage collection + DISABLE_DELETE?: string; + // Repositories that can be pulled without credentials, see src/anonymous.ts + ANONYMOUS_PULL_REPOSITORIES?: string; +} + +/** The env seen by the routes: the configuration plus state that fetch() sets for each request. */ +export interface Env extends RegistryEnv { REGISTRY_CLIENT: Registry; + ANONYMOUS_REQUEST?: boolean; } const router = Router(); @@ -35,8 +54,8 @@ router.all("/v2/*", v2Router.fetch); router.all("*", () => new Response("Not Found.", { status: 404 })); -export default { - async fetch(request: Request, env: Env, context?: ExecutionContext) { +const handler = { + async fetch(request: Request, env: RegistryEnv, context?: ExecutionContext): Promise { if (!ensureConfig(env)) { return new AuthErrorResponse(request); } @@ -46,16 +65,26 @@ export default { return new AuthErrorResponse(request); } + let anonymous = false; const credentials = await authMethod.checkCredentials(request); if (!credentials.verified) { - console.warn(`Not Authorized. authmode=${authMethod.authmode}. verified=false`); - return new AuthErrorResponse(request); + // Requests without any credentials may pull repositories listed in ANONYMOUS_PULL_REPOSITORIES. + // Wrong or expired credentials are still rejected, so clients notice them. /v2/ keeps answering + // 401 with a Basic challenge, which is what makes docker send credentials for pushes. + anonymous = request.headers.get("Authorization") === null && anonymousPullRepository(env, request) !== null; + if (!anonymous) { + console.warn(`Not Authorized. authmode=${authMethod.authmode}. verified=false`); + return new AuthErrorResponse(request); + } } - env.REGISTRY_CLIENT = new R2Registry(env); + // env is shared by all concurrent requests of this isolate, so everything that depends on the + // request goes into a copy. + const requestEnv = { ...env, ANONYMOUS_REQUEST: anonymous } as Env; + requestEnv.REGISTRY_CLIENT = new R2Registry(requestEnv); try { // Dispatch the request to the appropriate route - const res = await router.fetch(request, env, context); + const res = await router.fetch(request, requestEnv, context); return res; } catch (err) { if (err instanceof Response) { @@ -79,9 +108,12 @@ export default { return new InternalError(); } }, -} satisfies ExportedHandler; +} satisfies ExportedHandler; + +export { handler }; +export default handler; -const ensureConfig = (env: Env): boolean => { +const ensureConfig = (env: RegistryEnv): boolean => { if (!env.REGISTRY) { console.error( "env.REGISTRY is not setup. Please setup an R2 bucket and add the binding in your wrangler config file. Try 'npx wrangler --env production r2 bucket create r2-registry'", diff --git a/package.json b/package.json index 398dcd4..59b7580 100644 --- a/package.json +++ b/package.json @@ -4,6 +4,18 @@ "description": "An open-source R2 registry", "type": "module", "main": "index.ts", + "types": "index.ts", + "exports": { + ".": { + "types": "./index.ts", + "default": "./index.ts" + }, + "./package.json": "./package.json" + }, + "files": [ + "index.ts", + "src" + ], "scripts": { "deploy": "wrangler deploy --minify --env production", "dev:miniflare": "wrangler dev --env dev --port 9999 --live-reload", diff --git a/src/anonymous.ts b/src/anonymous.ts new file mode 100644 index 0000000..bc30bee --- /dev/null +++ b/src/anonymous.ts @@ -0,0 +1,55 @@ +import type { RegistryEnv } from ".."; +import { isValidDigest } from "./user"; + +// Read-only registry endpoints of a repository. Everything else (/v2/, /v2/_catalog, uploads, +// deletes, gc) always needs credentials. +const readPaths: { pattern: RegExp; digest?: boolean }[] = [ + { pattern: /^\/v2\/(.+)\/manifests\/[^/]+$/ }, + { pattern: /^\/v2\/(.+)\/blobs\/([^/]+)$/, digest: true }, + { pattern: /^\/v2\/(.+)\/tags\/list$/ }, + { pattern: /^\/v2\/(.+)\/referrers\/[^/]+$/ }, +]; + +// Turns ANONYMOUS_PULL_REPOSITORIES into regular expressions. +// +// The variable is a comma or whitespace separated list of repository names, `*` matches any +// characters including `/`. Examples: `*` (every repository), `some-org/*`, `some-org/app,other/tool`. +export function anonymousPullPatterns(env: RegistryEnv): RegExp[] { + return (env.ANONYMOUS_PULL_REPOSITORIES ?? "") + .split(/[\s,]+/) + .filter((p) => p.length > 0) + .map((p) => new RegExp(`^${p.split("*").map(escapeRegExp).join(".*")}$`)); +} + +function escapeRegExp(s: string): string { + return s.replace(/[.*+?^${}()|[\]\\]/g, "\\$&"); +} + +// Returns the repository name if the request is a pull that may be served without credentials. +export function anonymousPullRepository(env: RegistryEnv, request: Request): string | null { + if (request.method !== "GET" && request.method !== "HEAD") { + return null; + } + + const patterns = anonymousPullPatterns(env); + if (patterns.length === 0) { + return null; + } + + const path = new URL(request.url).pathname; + for (const { pattern, digest } of readPaths) { + const match = pattern.exec(path); + if (match === null) { + continue; + } + + const name = match[1]; + if (digest && !isValidDigest(match[2])) { + return null; + } + + return patterns.some((p) => p.test(name)) ? name : null; + } + + return null; +} diff --git a/src/authentication-method.ts b/src/authentication-method.ts index aee6906..9213a47 100644 --- a/src/authentication-method.ts +++ b/src/authentication-method.ts @@ -1,9 +1,9 @@ -import { Env } from ".."; +import type { RegistryEnv } from ".."; import { newRegistryTokens } from "./token"; import { UserAuthenticator } from "./user"; import type { AuthenticatorCredentials } from "./user"; -export async function authenticationMethodFromEnv(env: Env) { +export async function authenticationMethodFromEnv(env: RegistryEnv) { if (env.JWT_REGISTRY_TOKENS_PUBLIC_KEY) { return await newRegistryTokens(env.JWT_REGISTRY_TOKENS_PUBLIC_KEY); } else if ((env.USERNAME && env.PASSWORD) || (env.READONLY_USERNAME && env.READONLY_PASSWORD)) { diff --git a/src/chunk.ts b/src/chunk.ts index d7cfeb2..3d66c8a 100644 --- a/src/chunk.ts +++ b/src/chunk.ts @@ -44,23 +44,28 @@ export async function getChunkBlob(env: Env, chunk: Chunk): Promise */ export function limit(streamInput: ReadableStream, limitBytes: number): ReadableStream { if (streamInput instanceof FixedLengthStream) return streamInput; + + // R2 only accepts streams of a known length, so this has to stay a FixedLengthStream. const stream = new FixedLengthStream(limitBytes, {}); (async () => { - const w = stream.writable.getWriter(); - const r = streamInput.getReader(); + const reader = streamInput.getReader(); + const writer = stream.writable.getWriter(); let written = 0; - while (true) { - const { done, value } = await r.read(); - if (done) break; - await w.write(value); - written += value.length; - if (written >= limitBytes) break; + try { + while (written < limitBytes) { + const { done, value } = await reader.read(); + if (done) break; + const toWrite = value.length + written > limitBytes ? value.slice(0, limitBytes - written) : value; + await writer.write(toWrite); + written += toWrite.length; + } + reader.releaseLock(); + await writer.close(); + } catch (e) { + // Propagate the error to the consumer instead of leaving it waiting for bytes that never come + await writer.abort(e).catch(() => undefined); } - - r.releaseLock(); - w.releaseLock(); - await stream.writable.close(); })(); return stream.readable; diff --git a/src/errors.ts b/src/errors.ts index 95e655f..3afb14d 100644 --- a/src/errors.ts +++ b/src/errors.ts @@ -72,7 +72,8 @@ export class InternalError extends Response { export class ManifestError extends Response { constructor( - code: "MANIFEST_INVALID" | "BLOB_UNKNOWN" | "MANIFEST_UNVERIFIED" | "TAG_INVALID" | "NAME_INVALID", + code: + "MANIFEST_INVALID" | "BLOB_UNKNOWN" | "MANIFEST_UNVERIFIED" | "TAG_INVALID" | "NAME_INVALID" | "DIGEST_INVALID", message: string, detail: Record = {}, ) { @@ -94,6 +95,68 @@ export class ManifestError extends Response { } } +export class ImmutableTagError extends Response { + constructor(reference: string, action = "reassigned") { + super( + JSON.stringify({ + errors: [ + { + code: "DENIED", + message: `immutable tag ${reference} cannot be ${action}`, + detail: { reference }, + }, + ], + }), + { + status: 409, + headers: { + "content-type": "application/json;charset=UTF-8", + }, + }, + ); + } +} + +export class ImmutableBlobError extends Response { + constructor(digest: string) { + super( + JSON.stringify({ + errors: [ + { + code: "DENIED", + message: `blob ${digest} cannot be deleted while immutable tag policy is enabled`, + detail: { digest }, + }, + ], + }), + { + status: 409, + headers: { + "content-type": "application/json;charset=UTF-8", + }, + }, + ); + } +} + +// Returned for every delete when DISABLE_DELETE is set. The OCI distribution spec allows a registry +// to disable deletion and answer 405 Method Not Allowed. +export class DeletionDisabledError extends Response { + constructor() { + super( + JSON.stringify({ + errors: [{ code: "UNSUPPORTED", message: "deleting is disabled on this registry", detail: null }], + }), + { + status: 405, + headers: { + "content-type": "application/json;charset=UTF-8", + }, + }, + ); + } +} + export class ServerError extends Response { constructor(message: string, errorCode = 500) { super(JSON.stringify({ errors: [{ code: "SERVER_ERROR", message, detail: null }] }), { diff --git a/src/manifest.ts b/src/manifest.ts index 234fa59..3cc8cec 100644 --- a/src/manifest.ts +++ b/src/manifest.ts @@ -1,6 +1,16 @@ import { z } from "zod"; -const dockerManifestListContentType = "application/vnd.docker.distribution.manifest.list.v2+json"; +export const ociImageManifestContentType = "application/vnd.oci.image.manifest.v1+json"; +export const ociImageIndexContentType = "application/vnd.oci.image.index.v1+json"; +export const dockerImageManifestContentType = "application/vnd.docker.distribution.manifest.v2+json"; +export const dockerManifestListContentType = "application/vnd.docker.distribution.manifest.list.v2+json"; + +const manifestContentTypes: ReadonlySet = new Set([ + ociImageManifestContentType, + ociImageIndexContentType, + dockerImageManifestContentType, + dockerManifestListContentType, +]); const platformSchema = z.object({ "architecture": z.string(), @@ -80,3 +90,40 @@ export const manifestSchema = z ); export type ManifestSchema = z.infer; + +/** + * The top-level mediaType is OPTIONAL in the OCI image-spec ("SHOULD be used"), and clients + * such as Helm omit it. The Content-Type header of the PUT carries the same information, so + * fill it in from there before validating. + * + * Only the parsed object is changed. Manifest bytes are stored verbatim, so the stored object + * still hashes to the digest the client pushed it under. + */ +export function withInferredMediaType(manifestJSON: unknown, contentType: string): unknown { + if (typeof manifestJSON !== "object" || manifestJSON === null || Array.isArray(manifestJSON)) { + return manifestJSON; + } + + const manifest = manifestJSON as Record; + if (manifest.schemaVersion !== 2 || manifest.mediaType !== undefined) { + return manifestJSON; + } + + const mediaType = contentType.split(";")[0].trim(); + if (!manifestContentTypes.has(mediaType)) { + return manifestJSON; + } + + return { ...manifest, mediaType }; +} + +/** A failed union reports a bare "Invalid input" at the root; the detail is in the per-option issues. */ +export function manifestIssueMessage(error: z.ZodError): string { + const issue = error.issues[0]; + if (issue === undefined) return "invalid manifest"; + + const candidates = issue.code === "invalid_union" ? issue.errors.flat() : [issue]; + const best = candidates.find((candidate) => candidate.path.length > 0) ?? issue; + const path = best.path.length ? `${best.path.join(".")}: ` : ""; + return `${path}${best.message}`; +} diff --git a/src/registry/http.ts b/src/registry/http.ts index 2b3dda8..fea47ac 100644 --- a/src/registry/http.ts +++ b/src/registry/http.ts @@ -17,6 +17,7 @@ import { RegistryError, UploadId, UploadObject, + BlobRangeRequest, } from "./registry"; import { ociImageIndexContentType } from "./r2"; @@ -156,18 +157,50 @@ function ctxIntoHeaders(ctx: HTTPContext): Headers { return headers; } -function ctxIntoRequest(ctx: HTTPContext, url: URL, method: string, path: string, body?: BodyInit): Request { +function ctxIntoRequest( + ctx: HTTPContext, + url: URL, + method: string, + path: string, + body?: BodyInit, + extraHeaders?: HeadersInit, +): Request { const urlReq = `${url.protocol}//${url.host}/v2${ ctx.repository === "" || ctx.repository === "/" ? "/" : ctx.repository + "/" }${path}`; + const headers = ctxIntoHeaders(ctx); + if (extraHeaders !== undefined) { + new Headers(extraHeaders).forEach((value, key) => headers.set(key, value)); + } return new Request(urlReq, { method, body, redirect: "follow", - headers: ctxIntoHeaders(ctx), + headers, }); } +// Serializes a blob range request into an HTTP "Range" request header value. +function rangeRequestHeader(range: BlobRangeRequest): string { + if ("suffix" in range) { + return `bytes=-${range.suffix}`; + } + + return `bytes=${range.offset}-${range.end === undefined ? "" : range.end}`; +} + +// Parses an HTTP "Content-Range: bytes -/" response header. +function parseContentRange(header: string | null): { start: number; end: number; size: number } | null { + if (header === null) return null; + const match = /^bytes (\d+)-(\d+)\/(\d+)$/.exec(header.trim()); + if (match === null) return null; + const start = Number(match[1]); + const end = Number(match[2]); + const size = Number(match[3]); + if (!Number.isInteger(start) || !Number.isInteger(end) || !Number.isInteger(size)) return null; + return { start, end, size }; +} + function authHeaderIntoAuthContext(urlObject: URL, authenticateHeader: string): AuthContext { const url = urlObject.toString(); const parts = authenticateHeader.split(" "); @@ -519,18 +552,28 @@ export class RegistryHTTPClient implements Registry { } } - async getLayer(name: string, digest: string): Promise { + async getLayer(name: string, digest: string, range?: BlobRangeRequest): Promise { const namespace = name.includes("/") || !isDockerDotIO(this.url) ? name : `library/${name}`; try { const ctx = await this.authenticate(namespace); - const req = ctxIntoRequest(ctx, this.url, "GET", `${namespace}/blobs/${digest}`); + const rangeHeader = range === undefined ? undefined : rangeRequestHeader(range); + const req = ctxIntoRequest( + ctx, + this.url, + "GET", + `${namespace}/blobs/${digest}`, + undefined, + rangeHeader !== undefined ? { Range: rangeHeader } : undefined, + ); let res = await fetch(req); if (!res.ok) { // This means we got a redirect, so let's try again this URL but // without any headers. Services like S3 reject authorization headers altogether // if the authentication is included in the URL. if (res.url !== req.url) { - const redirectResponse = await fetch(new Request(res.url)); + const redirectResponse = await fetch( + new Request(res.url, rangeHeader !== undefined ? { headers: { Range: rangeHeader } } : undefined), + ); if (!redirectResponse.ok) { return { response: res, @@ -549,11 +592,27 @@ export class RegistryHTTPClient implements Registry { throw new Error("returned body is null"); } - return { + const layer: GetLayerResponse = { stream: res.body, size: +(res.headers.get("Content-Length") ?? "0"), digest: res.headers.get("Digest-Content-Digest") ?? digest, }; + + // If we asked for a range and the upstream honored it, surface the partial-content metadata so + // the caller can reply with 206. A 200 here means the upstream ignored the range and we serve + // the full blob (best-effort). Serving a partial body as if it were complete would corrupt it, + // so a 206 without a parseable Content-Range is treated as an error. + if (range !== undefined && res.status === 206) { + const contentRange = parseContentRange(res.headers.get("Content-Range")); + if (contentRange === null) { + throw new Error("upstream returned 206 without a parseable Content-Range header"); + } + + layer.size = contentRange.size; + layer.contentRange = contentRange; + } + + return layer; } catch (err) { console.error(`Error doing get layer with ${namespace} and ${digest}: ` + errorString(err)); return { diff --git a/src/registry/r2.ts b/src/registry/r2.ts index 7f01098..c80138c 100644 --- a/src/registry/r2.ts +++ b/src/registry/r2.ts @@ -9,10 +9,10 @@ import { limit, split, } from "../chunk"; -import { InternalError, ManifestError, RangeError, ServerError } from "../errors"; +import { ImmutableTagError, InternalError, ManifestError, RangeError, ServerError } from "../errors"; import { SHA256_PREFIX_LEN, getSHA256, hexToDigest, isValidDigest } from "../user"; -import { readableToBlob, readerToBlob, wrap } from "../utils"; -import { BlobUnknownError, ManifestUnknownError } from "../v2-errors"; +import { errorString, jsonHeaders, readableToBlob, readerToBlob, wrap } from "../utils"; +import { BlobUnknownError, DigestInvalidError, ManifestUnknownError } from "../v2-errors"; import { CheckLayerResponse, CheckManifestResponse, @@ -28,11 +28,58 @@ import { UploadId, UploadObject, wrapError, + BlobRangeRequest, } from "./registry"; import { GarbageCollectionMode, GarbageCollector } from "./garbage-collector"; -import { ManifestSchema, manifestSchema } from "../manifest"; +import { ManifestSchema, manifestSchema, manifestIssueMessage, withInferredMediaType } from "../manifest"; +import { isImmutableTagReference, resolveImmutableTagPattern } from "./tag-policy"; -export const ociImageIndexContentType = "application/vnd.oci.image.index.v1+json"; +export { ociImageIndexContentType } from "../manifest"; + +function rangeNotSatisfiableResponse(size: number): Response { + return new Response(null, { + status: 416, + headers: { + "Content-Range": `bytes */${size}`, + "Accept-Ranges": "bytes", + }, + }); +} + +/** + * Writes a content-addressed object only if it does not exist yet, and treats an existing object as + * success. These keys (a blob or manifest stored under its digest, a mounted blob, a referrer entry) + * always describe the same content, so an existing object means the write already happened. Never + * overwriting them keeps pushes working on buckets whose content prefixes refuse overwrites, such as + * R2 buckets with bucket locks or other retention rules, and avoids rewriting large blobs on re-push. + * + * Returns the stored object, or the object that was already there. + */ +export async function putIfAbsent( + bucket: R2Bucket, + key: string, + value: ReadableStream | ArrayBuffer | ArrayBufferView | string | null | Blob, + options: R2PutOptions = {}, +): Promise { + let created: R2Object | null; + try { + created = await bucket.put(key, value, { ...options, onlyIf: { etagDoesNotMatch: "*" } }); + } catch (err) { + // A retention rule can refuse the write outright instead of failing the precondition. That is + // still success when the object is already there. + const existing = await bucket.head(key); + if (existing !== null) return existing; + throw err; + } + + if (created !== null) return created; + const existing = await bucket.head(key); + if (existing === null) { + throw new Error(`conditional write of ${key} was refused, but the object does not exist`); + } + + return existing; +} function referrersPrefix(name: string, digest: string): string { return `${name}/_referrers/${digest}/`; @@ -198,6 +245,8 @@ export async function encodeState(state: State, env: Env): Promise<{ jwt: string } export const symlinkHeader = "X-Serverless-Registry-Symlink"; +export const symlinkDigestHeader = "X-Serverless-Registry-Symlink-Digest"; +export const symlinkSizeHeader = "X-Serverless-Registry-Symlink-Size"; export async function getUploadState( name: string, @@ -359,15 +408,10 @@ export class R2Registry implements Registry { async verifyManifest(name: string, manifest: ManifestSchema) { if (manifest.schemaVersion === 2 && "manifests" in manifest) { - for (const manifestElement of manifest.manifests) { - const key = manifestElement.digest; - const res = await this.env.REGISTRY.head(`${name}/manifests/${key}`); - if (res === null) { - console.error(`Manifest with digest ${key} doesn't exist`); - return new ManifestError("BLOB_UNKNOWN", `unknown manifest ${key}`); - } - } - + // Sparse indexes are allowed: when only some platforms of an image are mirrored (regsync + // `platforms`), the original index is pushed unchanged so its digest stays the one upstream + // has, and the manifests of the other platforms simply do not exist here. Pulling one of + // those platforms fails with MANIFEST_UNKNOWN. return null; } @@ -509,6 +553,17 @@ export class R2Registry implements Registry { shaWriter.close(); const digest = await sha256.digest; const digestStr = hexToDigest(digest); + + // A reference containing ":" addresses the manifest by digest (an OCI tag never contains ":"). + // The submitted content must hash to exactly that digest; a mismatched or malformed digest + // reference is a client error (400 DIGEST_INVALID), not a tag to store under the wrong key. + if (reference.includes(":") && reference !== digestStr) { + const message = isValidDigest(reference) + ? `provided digest ${reference} does not match content digest ${digestStr}` + : `invalid digest reference ${reference}`; + return { response: new ManifestError("DIGEST_INVALID", message) }; + } + const text = await blob.text(); let manifestJSON: unknown; try { @@ -519,12 +574,10 @@ export class R2Registry implements Registry { }; } - const manifestResult = manifestSchema.safeParse(manifestJSON); + const manifestResult = manifestSchema.safeParse(withInferredMediaType(manifestJSON, contentType)); if (!manifestResult.success) { - const firstIssue = manifestResult.error.issues[0]; - const path = firstIssue?.path.length ? `${firstIssue.path.join(".")}: ` : ""; return { - response: new ManifestError("MANIFEST_INVALID", `${path}${firstIssue?.message ?? "invalid manifest"}`), + response: new ManifestError("MANIFEST_INVALID", manifestIssueMessage(manifestResult.error)), }; } @@ -535,17 +588,9 @@ export class R2Registry implements Registry { response: new ManifestError("MANIFEST_INVALID", `invalid subject digest ${subjectDigest}`), }; } - if (subjectDigest !== undefined) { - const [subjectManifest, subjectManifestErr] = await wrap(env.REGISTRY.head(`${name}/manifests/${subjectDigest}`)); - if (subjectManifestErr) { - return wrapError("putManifestInner", subjectManifestErr); - } - if (subjectManifest === null) { - return { - response: new ManifestError("BLOB_UNKNOWN", `unknown subject ${subjectDigest}`), - }; - } - } + // The subject does not have to exist yet: the OCI distribution spec (v1.1) allows pushing a + // referrer before its subject, and copy tools like regsync push the manifests of an index in + // parallel, so buildx attestations regularly arrive before the image they describe. const referrerDescriptor = descriptorFromManifest(manifest, digestStr, blob.size); if (checkLayers) { @@ -562,42 +607,66 @@ export class R2Registry implements Registry { hasSubject: subjectDigest !== undefined ? "true" : "false", ...(subjectDigest !== undefined ? { subjectDigest } : {}), }; + const immutablePattern = resolveImmutableTagPattern(env.IMMUTABLE_TAG_PATTERN); - const putReference = async () => { - // if the reference is the same as a digest, it's not necessary to insert - if (reference === digestStr) return; - return await env.REGISTRY.put(`${name}/manifests/${reference}`, text, { - sha256: digest, - httpMetadata: { - contentType, - }, - customMetadata, - }); + const putOptions = { + sha256: digest, + httpMetadata: { + contentType, + }, + customMetadata, }; - const putTasks: Promise[] = [ - putReference(), - // this is the "main" manifest - env.REGISTRY.put(`${name}/manifests/${digestStr}`, text, { - sha256: digest, - httpMetadata: { - contentType, - }, - customMetadata, - }), - ]; - - if (referrerDescriptor !== null && subjectDigest !== undefined) { - putTasks.push( - env.REGISTRY.put(referrersPath(name, subjectDigest, digestStr), JSON.stringify(referrerDescriptor), { + const putReferrer = () => { + if (referrerDescriptor === null || subjectDigest === undefined) { + return null; + } + return putIfAbsent( + env.REGISTRY, + referrersPath(name, subjectDigest, digestStr), + JSON.stringify(referrerDescriptor), + { httpMetadata: { contentType: "application/json", }, - }), + }, ); - } + }; + const digestPut = () => putIfAbsent(env.REGISTRY, `${name}/manifests/${digestStr}`, text, putOptions); + const immutableReference = reference !== digestStr && isImmutableTagReference(reference, immutablePattern); + + if (immutableReference) { + // The digest must be durable before the protected tag becomes visible, + // and only an accepted tag may publish derived referrer state. + await digestPut(); + const referenceKey = `${name}/manifests/${reference}`; + const created = await env.REGISTRY.put(referenceKey, text, { + ...putOptions, + onlyIf: { etagDoesNotMatch: "*" }, + }); + if (created === null) { + const existing = await env.REGISTRY.head(referenceKey); + const existingDigest = existing?.checksums.sha256 ? hexToDigest(existing.checksums.sha256) : undefined; + if (existingDigest !== digestStr) { + return { response: new ImmutableTagError(reference) }; + } + } - await Promise.all(putTasks); + const referrerPut = putReferrer(); + if (referrerPut !== null) { + await referrerPut; + } + } else { + const putTasks: Promise[] = [digestPut()]; + if (reference !== digestStr) { + putTasks.push(env.REGISTRY.put(`${name}/manifests/${reference}`, text, putOptions)); + } + const referrerPut = putReferrer(); + if (referrerPut !== null) { + putTasks.push(referrerPut); + } + await Promise.all(putTasks); + } return { digest: hexToDigest(digest), location: `/v2/${name}/manifests/${reference}`, @@ -650,11 +719,22 @@ export class R2Registry implements Registry { // Trying to mount a layer from sourceLayerPath to destinationLayerPath // Create linked file with custom metadata + if (res.checksums.sha256 === null) { + return { response: new ServerError("invalid checksum from R2 backend") }; + } const [newFile, error] = await wrap( - this.env.REGISTRY.put(destinationLayerPath, sourceLayerPath, { + putIfAbsent(this.env.REGISTRY, destinationLayerPath, sourceLayerPath, { + // Symlink object content is the source blob path string. + // The object checksum must match the symlink payload to satisfy R2. sha256: await getSHA256(sourceLayerPath, ""), httpMetadata: res.httpMetadata, - customMetadata: { [symlinkHeader]: sourceName }, // Storing target repository name in metadata (to easily resolve recursive layer mounting) + customMetadata: { + // Storing target repository name in metadata (to easily resolve recursive layer mounting) + [symlinkHeader]: sourceName, + // Store source layer metadata so HEAD can answer without loading symlink body. + [symlinkDigestHeader]: digest, + [symlinkSizeHeader]: `${res.size}`, + }, }), ); if (error) { @@ -683,6 +763,63 @@ export class R2Registry implements Registry { }; } + const expectedDigest = tag.startsWith("sha256:") ? tag : null; + const actualDigest = res.checksums.sha256 ? hexToDigest(res.checksums.sha256) : null; + const symlinkByChecksumMismatch = + expectedDigest !== null && actualDigest !== null && actualDigest !== expectedDigest; + const symlinkMetadata = res.customMetadata ?? {}; + const symlinkByMetadata = symlinkHeader in symlinkMetadata; + + // Handle R2 symlink layers. + // We detect symlinks by: + // 1) explicit metadata, or + // 2) checksum mismatch between requested digest and R2 object checksum + // (the symlink object checksum is based on symlink payload, not mounted blob bytes). + if (symlinkByMetadata || symlinkByChecksumMismatch) { + // Fast path for symlinks created by newer versions that include source metadata. + const metadataSize = +(symlinkMetadata[symlinkSizeHeader] ?? ""); + const metadataDigest = symlinkMetadata[symlinkDigestHeader]; + if (Number.isFinite(metadataSize) && metadataSize >= 0 && metadataDigest) { + return { + digest: metadataDigest, + size: metadataSize, + exists: true, + }; + } + + // Backward-compatibility path for old symlinks: resolve link body and query target. + const [obj, getErr] = await wrap(this.env.REGISTRY.get(`${name}/blobs/${tag}`)); + if (getErr) { + return wrapError("layerExists", getErr); + } + if (!obj) { + return { exists: false }; + } + + const layerPath = await obj.text(); + const [linkName, linkDigest] = layerPath.split("/blobs/"); + if (!linkName || !linkDigest) { + // Backward compatibility: if this does not look like a symlink payload, + // fall back to object metadata from HEAD. + if (res.checksums.sha256 === null) { + return { response: new ServerError("invalid checksum from R2 backend") }; + } + + return { + digest: hexToDigest(res.checksums.sha256!), + size: res.size, + exists: true, + }; + } + + // Prevent recursive self-reference. + if (linkName === name && linkDigest === tag) { + return { exists: false }; + } + + return await this.env.REGISTRY_CLIENT.layerExists(linkName, linkDigest); + } + return { digest: hexToDigest(res.checksums.sha256!), size: res.size, @@ -690,38 +827,122 @@ export class R2Registry implements Registry { }; } - async getLayer(name: string, digest: string): Promise { - const [res, err] = await wrap(this.env.REGISTRY.get(`${name}/blobs/${digest}`)); - if (err) { - return wrapError("getLayer", err); + async getLayer(name: string, digest: string, range?: BlobRangeRequest): Promise { + const key = `${name}/blobs/${digest}`; + + if (range === undefined) { + const [res, err] = await wrap(this.env.REGISTRY.get(key)); + if (err) { + return wrapError("getLayer", err); + } + + if (!res) { + return { + response: new Response(JSON.stringify(BlobUnknownError), { status: 404 }), + }; + } + + // Handle R2 symlink + if (res.customMetadata && symlinkHeader in res.customMetadata) { + return await this.followLayerSymlink(name, digest, res, undefined); + } + + return { + stream: res.body!, + digest: hexToDigest(res.checksums.sha256!), + size: res.size, + }; } - if (!res) { + // Ranged read: inspect object metadata first so we can validate the requested range and + // resolve symlinks without streaming the full object. + const [head, headErr] = await wrap(this.env.REGISTRY.head(key)); + if (headErr) { + return wrapError("getLayer", headErr); + } + + if (!head) { return { response: new Response(JSON.stringify(BlobUnknownError), { status: 404 }), }; } - // Handle R2 symlink - if (res.customMetadata && symlinkHeader in res.customMetadata) { - const layerPath = await res.text(); - // Symlink detected! Will download layer from "layerPath" - const [linkName, linkDigest] = layerPath.split("/blobs/"); - if (linkName == name && linkDigest == digest) { + if (head.customMetadata && symlinkHeader in head.customMetadata) { + const [link, linkErr] = await wrap(this.env.REGISTRY.get(key)); + if (linkErr) { + return wrapError("getLayer", linkErr); + } + + if (!link) { return { response: new Response(JSON.stringify(BlobUnknownError), { status: 404 }), }; } - return await this.env.REGISTRY_CLIENT.getLayer(linkName, linkDigest); + + return await this.followLayerSymlink(name, digest, link, range); + } + + const totalSize = head.size; + let start: number; + let end: number; + if ("suffix" in range) { + // A suffix longer than the object is satisfied by the whole object, but a zero-length suffix + // selects no bytes at all and cannot be satisfied. + if (range.suffix <= 0 || totalSize === 0) { + return { response: rangeNotSatisfiableResponse(totalSize) }; + } + + start = Math.max(totalSize - range.suffix, 0); + end = totalSize - 1; + } else { + start = range.offset; + if (start < 0 || start >= totalSize) { + return { response: rangeNotSatisfiableResponse(totalSize) }; + } + + end = range.end === undefined ? totalSize - 1 : Math.min(range.end, totalSize - 1); + if (end < start) { + return { response: rangeNotSatisfiableResponse(totalSize) }; + } + } + + const length = end - start + 1; + const [res, err] = await wrap(this.env.REGISTRY.get(key, { range: { offset: start, length } })); + if (err) { + return wrapError("getLayer", err); + } + + if (!res) { + return { + response: new Response(JSON.stringify(BlobUnknownError), { status: 404 }), + }; } return { stream: res.body!, - digest: hexToDigest(res.checksums.sha256!), - size: res.size, + digest: hexToDigest(head.checksums.sha256!), + size: totalSize, + contentRange: { start, end, size: totalSize }, }; } + private async followLayerSymlink( + name: string, + digest: string, + object: R2ObjectBody, + range: BlobRangeRequest | undefined, + ): Promise { + const layerPath = await object.text(); + // Symlink detected! Will download layer from "layerPath" + const [linkName, linkDigest] = layerPath.split("/blobs/"); + if (linkName == name && linkDigest == digest) { + return { + response: new Response(JSON.stringify(BlobUnknownError), { status: 404 }), + }; + } + return await this.env.REGISTRY_CLIENT.getLayer(linkName, linkDigest, range); + } + async startUpload(namespace: string): Promise { // Generate a unique ID for this upload const uuid = crypto.randomUUID(); @@ -898,8 +1119,10 @@ export class R2Registry implements Registry { httpMetadata: new Headers(headers), customMetadata: headers, }); - state.parts.push(await partTask); - await r2RegistryObjectTask; + + // Run both in parallel and wait for both to complete + const [part] = await Promise.all([partTask, r2RegistryObjectTask]); + state.parts.push(part); return; } @@ -927,7 +1150,9 @@ export class R2Registry implements Registry { }; if (length === undefined) { - console.error("Length needs to be defined"); + console.error( + "Length needs to be defined for streaming upload to R2. Ensure Content-Length or Content-Range is provided.", + ); return { response: new InternalError(), }; @@ -975,14 +1200,36 @@ export class R2Registry implements Registry { const state = hashedState.state; const uuid = state.registryUploadId; - if (state.parts.length === 0) { - if (!stream) { - console.error("There has been an upload with zero parts and the body is null"); + + // Commit the finished content under the client-claimed digest. R2 verifies the sha256 we + // hand it against the bytes it stored, so a digest the client got wrong surfaces here as a + // checksum-mismatch rejection — translate that into a 400 DIGEST_INVALID rather than letting + // it bubble up as an opaque 500. Any other failure is a genuine server error. + const putBlob = async (body: ReadableStream | Uint8Array | null): Promise => { + const [, err] = await wrap( + putIfAbsent(this.env.REGISTRY, `${namespace}/blobs/${expectedSha}`, body, { + sha256: (expectedSha as string).slice(SHA256_PREFIX_LEN), + }), + ); + if (err === null) return null; + const message = errorString(err); + // Matches the wording R2 uses when the stored bytes don't hash to the requested sha256. + // This couples to the runtime's error text; the unit test asserts the 400 body so a future + // wording change surfaces as a test failure rather than a silent regression to 500. + if (/checksum|did not match/i.test(message)) { return { - response: new InternalError(), + response: new Response(JSON.stringify(DigestInvalidError()), { status: 400, headers: jsonHeaders() }), }; } + console.error("finishUpload put failed:", message); + return { response: new InternalError() }; + }; + if (state.parts.length === 0) { + // No multipart parts were staged: the whole blob arrives in this request body (a monolithic + // PUT), or it is a zero-byte blob. An absent body is a valid empty blob — store empty bytes + // (R2 requires a known-length body, so a length-less empty stream cannot be used here) and + // let the checksum check confirm the client really claimed the empty digest. if (length && length > MAXIMUM_CHUNK) { console.error("Surpasses MAXIMUM_CHUNK"); return { @@ -990,20 +1237,19 @@ export class R2Registry implements Registry { }; } - await this.env.REGISTRY.put(`${namespace}/blobs/${expectedSha}`, stream, { - sha256: (expectedSha as string).slice(SHA256_PREFIX_LEN), - }); + // With bytes to store, the request body carries its own (known) length. With none, it is a + // zero-byte blob — hand R2 empty bytes rather than a length-less empty stream, which it rejects. + const putErr = await putBlob(length && length > 0 ? stream! : new Uint8Array(0)); + if (putErr) return putErr; } else { const upload = this.env.REGISTRY.resumeMultipartUpload(uuid, state.uploadId); - // TODO: Handle one last buffer here + // A final chunk carried by the finalizing PUT is appended beforehand via uploadChunk (the + // same path a PATCH uses), so the staged parts are complete here. See the PUT handler. await upload.complete(state.parts); const obj = await this.env.REGISTRY.get(uuid); - const put = this.env.REGISTRY.put(`${namespace}/blobs/${expectedSha}`, obj!.body, { - sha256: (expectedSha as string).slice(SHA256_PREFIX_LEN), - }); - - await put; + const putErr = await putBlob(obj!.body); await this.env.REGISTRY.delete(uuid); + if (putErr) return putErr; } await this.env.REGISTRY.delete(getRegistryUploadsPath(state)); @@ -1048,7 +1294,7 @@ export class R2Registry implements Registry { return false; } - await this.env.REGISTRY.put(`${namespace}/blobs/${sha256}`, stream, { + await putIfAbsent(this.env.REGISTRY, `${namespace}/blobs/${sha256}`, stream, { sha256: (sha256 as string).slice(SHA256_PREFIX_LEN), }); return { diff --git a/src/registry/registry.ts b/src/registry/registry.ts index 217ed7e..e3f5416 100644 --- a/src/registry/registry.ts +++ b/src/registry/registry.ts @@ -25,6 +25,11 @@ const registryConfiguration = z export type RegistryConfiguration = z.infer; export function registries(env: Env): RegistryConfiguration[] { + // Anonymous clients must not be able to make us copy arbitrary upstream images into R2 + if (env.ANONYMOUS_REQUEST) { + return []; + } + if (env.REGISTRIES_JSON === undefined || env.REGISTRIES_JSON.length === 0) { return []; } @@ -99,11 +104,17 @@ export type GetManifestResponse = { contentType: string; }; +// requested byte range for a layer read. Either an offset with an inclusive end (end omitted means +// until the end of the object), or a suffix asking for the last N bytes of the object. +export type BlobRangeRequest = { offset: number; end?: number } | { suffix: number }; + // returned by getLayer when it successfully retrieves a layer export type GetLayerResponse = { stream: ReadableStream; digest: string; size: number; + // present when the response is a partial (ranged) read, drives the 206 Partial Content response + contentRange?: { start: number; end: number; size: number }; }; export type ReferrerDescriptor = { @@ -150,7 +161,7 @@ export interface Registry { layerExists(namespace: string, digest: string): Promise; // get a layer stream from the registry - getLayer(namespace: string, digest: string): Promise; + getLayer(namespace: string, digest: string, range?: BlobRangeRequest): Promise; // list referrers for a subject digest listReferrers( diff --git a/src/registry/tag-policy.ts b/src/registry/tag-policy.ts new file mode 100644 index 0000000..a2cbda1 --- /dev/null +++ b/src/registry/tag-policy.ts @@ -0,0 +1,17 @@ +import { isValidDigest } from "../user"; + +export function resolveImmutableTagPattern(source: string | undefined): RegExp | null { + const trimmed = source?.trim(); + if (!trimmed) return null; + + try { + return new RegExp(`^(?:${trimmed})$`); + } catch (error) { + const detail = error instanceof Error ? error.message : String(error); + throw new Error(`Invalid IMMUTABLE_TAG_PATTERN: ${detail}`, { cause: error }); + } +} + +export function isImmutableTagReference(reference: string, pattern: RegExp | null): boolean { + return pattern !== null && !isValidDigest(reference) && pattern.test(reference); +} diff --git a/src/router.ts b/src/router.ts index f8de37a..de1cbf0 100644 --- a/src/router.ts +++ b/src/router.ts @@ -1,12 +1,13 @@ import { Router } from "itty-router"; import { BlobUnknownError, ManifestUnknownError } from "./v2-errors"; -import { InternalError, ServerError } from "./errors"; -import { errorString, jsonHeaders, wrap } from "./utils"; +import { DeletionDisabledError, ImmutableBlobError, ImmutableTagError, InternalError, ServerError } from "./errors"; +import { errorString, getStreamSize, jsonHeaders, wrap } from "./utils"; import { hexToDigest, isValidDigest } from "./user"; import { ManifestTagsListTooBigError } from "./v2-responses"; import { Env } from ".."; import { MINIMUM_CHUNK, MAXIMUM_CHUNK, MAXIMUM_CHUNK_UPLOAD_SIZE } from "./chunk"; import { + BlobRangeRequest, CheckLayerResponse, CheckManifestResponse, FinishedUploadObject, @@ -19,6 +20,7 @@ import { } from "./registry/registry"; import { RegistryHTTPClient } from "./registry/http"; import { ociImageIndexContentType } from "./registry/r2"; +import { isImmutableTagReference, resolveImmutableTagPattern } from "./registry/tag-policy"; const maxReferrersListLimit = 1000; const isOpaqueReferrersCursor = (cursor: string) => cursor.startsWith("/v2/"); @@ -27,8 +29,19 @@ function formatNextLink(url: URL): string { return `<${url.toString()}>; rel="next"`; } +// Stops the runtime's automatic gzip, which switches the response to chunked transfer-encoding and +// drops the Content-Length the distribution spec requires on blob/manifest GET and HEAD. Bodies are +// served verbatim. Spread into each such response's headers. +const identityEncoding = { "Content-Encoding": "identity" } as const; + const v2Router = Router({ base: "/v2/" }); +// DISABLE_DELETE turns the registry into an append-only store: manifests, blobs and garbage +// collection can't delete anything. Overwriting a mutable tag is still allowed. +export function deletionDisabled(env: Env): boolean { + return ["true", "1", "yes"].includes((env.DISABLE_DELETE ?? "").trim().toLowerCase()); +} + v2Router.get("/", async (_req, _env: Env) => { return new Response(); }); @@ -73,8 +86,19 @@ v2Router.delete("/:name+/manifests/:reference", async (req, env: Env) => { // // If somehow we need to remove by paginating, we accept a last query param. + if (deletionDisabled(env)) { + return new DeletionDisabledError(); + } + const { last, limit } = req.query; const { name, reference } = req.params; + const immutablePattern = resolveImmutableTagPattern(env.IMMUTABLE_TAG_PATTERN); + if (immutablePattern !== null && isValidDigest(reference)) { + return new ImmutableTagError(reference, "deleted while immutable tag policy is enabled"); + } + if (isImmutableTagReference(reference, immutablePattern)) { + return new ImmutableTagError(reference, "deleted"); + } const manifest = await env.REGISTRY.head(`${name}/manifests/${reference}`); if (manifest === null) { return new Response(JSON.stringify(ManifestUnknownError(reference)), { status: 404, headers: jsonHeaders() }); @@ -101,15 +125,23 @@ v2Router.delete("/:name+/manifests/:reference", async (req, env: Env) => { limit: limitInt, cursor: last?.toString(), }); + const aliasesToDelete: string[] = []; for (const tag of tags.objects) { if (!tag.checksums.sha256) { continue; } if (hexToDigest(tag.checksums.sha256) === reference && tag.key !== `${name}/manifests/${reference}`) { - await env.REGISTRY.delete(tag.key); + const tagReference = tag.key.slice(`${name}/manifests/`.length); + if (isImmutableTagReference(tagReference, immutablePattern)) { + return new ImmutableTagError(tagReference, "deleted through its manifest digest"); + } + aliasesToDelete.push(tag.key); } } + if (aliasesToDelete.length > 0) { + await env.REGISTRY.delete(aliasesToDelete); + } const url = new URL(req.url); if (tags.truncated) { @@ -148,6 +180,7 @@ v2Router.head("/:name+/manifests/:reference", async (req, env: Env) => { "Content-Length": res.size.toString(), "Content-Type": res.contentType, "Docker-Content-Digest": res.digest, + ...identityEncoding, }, }); } @@ -207,6 +240,7 @@ v2Router.head("/:name+/manifests/:reference", async (req, env: Env) => { "Content-Length": checkManifestResponse.size.toString(), "Content-Type": checkManifestResponse.contentType, "Docker-Content-Digest": checkManifestResponse.digest, + ...identityEncoding, }, }); }); @@ -220,6 +254,7 @@ v2Router.get("/:name+/manifests/:reference", async (req, env: Env, context: Exec "Content-Length": res.size.toString(), "Content-Type": res.contentType, "Docker-Content-Digest": res.digest, + ...identityEncoding, }, }); } @@ -270,6 +305,7 @@ v2Router.get("/:name+/manifests/:reference", async (req, env: Env, context: Exec "Content-Length": getManifestResponse.size.toString(), "Content-Type": getManifestResponse.contentType, "Docker-Content-Digest": getManifestResponse.digest, + ...identityEncoding, }, }); }); @@ -358,54 +394,105 @@ v2Router.get("/:name+/referrers/:digest", async (req, env: Env) => { ); }); +// Parses a single HTTP byte range request of the form "bytes=-", "bytes=-" or +// the suffix form "bytes=-", which asks for the last n bytes. +// Multi-range and malformed values are ignored so the full object is served. +function parseBlobRange(header: string | null): BlobRangeRequest | undefined { + if (header === null) return undefined; + const match = /^bytes=(\d*)-(\d*)$/.exec(header.trim()); + if (match === null) return undefined; + const [, startValue, endValue] = match; + if (startValue === "") { + // "bytes=-" has neither a start nor a suffix length, so there is nothing to satisfy. + if (endValue === "") return undefined; + const suffix = Number(endValue); + return Number.isInteger(suffix) ? { suffix } : undefined; + } + + const offset = Number(startValue); + if (!Number.isInteger(offset)) return undefined; + if (endValue === "") return { offset }; + const end = Number(endValue); + if (!Number.isInteger(end)) return { offset }; + return { offset, end }; +} + +function blobGetResponse(layer: GetLayerResponse): Response { + const headers: Record = { + "Docker-Content-Digest": layer.digest, + "Accept-Ranges": "bytes", + ...identityEncoding, + }; + if (layer.contentRange !== undefined) { + const { start, end, size } = layer.contentRange; + headers["Content-Length"] = `${end - start + 1}`; + headers["Content-Range"] = `bytes ${start}-${end}/${size}`; + return new Response(layer.stream, { status: 206, headers }); + } + + headers["Content-Length"] = `${layer.size}`; + return new Response(layer.stream, { headers }); +} + v2Router.get("/:name+/blobs/:digest", async (req, env: Env, context: ExecutionContext) => { const { name, digest } = req.params; - const res = await env.REGISTRY_CLIENT.getLayer(name, digest); + const range = parseBlobRange(req.headers.get("range")); + const res = await env.REGISTRY_CLIENT.getLayer(name, digest, range); if (!("response" in res)) { - return new Response(res.stream, { - headers: { - "Docker-Content-Digest": res.digest, - "Content-Length": `${res.size}`, - }, - }); + return blobGetResponse(res); + } + + // A requested range that cannot be satisfied is reported directly instead of falling back to + // other registries. + if (res.response.status === 416) { + return res.response; } let layerResponse: GetLayerResponse | null = null; const registriesList = registries(env); for (const registry of registriesList) { const client = new RegistryHTTPClient(env, registry); - const response = await client.getLayer(name, digest); + const response = await client.getLayer(name, digest, range); if ("response" in response) { + // The blob exists upstream but the requested range doesn't fit it. Blobs are content + // addressed, so every registry holding this digest holds the same bytes and would answer the + // same way. Report it instead of letting it fall through to the 404 below, which would tell + // the client the blob doesn't exist and hide the object size it needs to retry. + if (response.response.status === 416) { + return response.response; + } + continue; } layerResponse = response; - const [s1, s2] = layerResponse.stream.tee(); - layerResponse.stream = s1; - context.waitUntil( - (async () => { - const [response, err] = await wrap(env.REGISTRY_CLIENT.monolithicUpload(name, digest, s2, layerResponse.size)); - if (err) { - console.error("Error uploading asynchronously the layer ", digest, "into main registry"); - return; - } + // Only cache full-object responses. A ranged/partial upstream response must never be written to + // R2 as if it were the complete blob, or the cached object would be corrupt. + if (range === undefined && layerResponse.contentRange === undefined) { + const fullLayer = layerResponse; + const [s1, s2] = fullLayer.stream.tee(); + fullLayer.stream = s1; + context.waitUntil( + (async () => { + const [response, err] = await wrap(env.REGISTRY_CLIENT.monolithicUpload(name, digest, s2, fullLayer.size)); + if (err) { + console.error("Error uploading asynchronously the layer ", digest, "into main registry"); + return; + } + + if (response === false) { + console.error("Layer might be too big for the registry client", fullLayer.size); + } + })(), + ); + } - if (response === false) { - console.error("Layer might be too big for the registry client", layerResponse.size); - } - })(), - ); break; } if (layerResponse === null) return new Response(JSON.stringify(BlobUnknownError), { status: 404 }); - return new Response(layerResponse.stream, { - headers: { - "Docker-Content-Digest": layerResponse.digest, - "Content-Length": `${layerResponse.size}`, - }, - }); + return blobGetResponse(layerResponse); }); v2Router.delete("/:name+/blobs/uploads/:id", async (req, env: Env) => { @@ -507,17 +594,19 @@ v2Router.get("/:name+/blobs/uploads/:uuid", async (req, env: Env) => { v2Router.patch("/:name+/blobs/uploads/:uuid", async (req, env: Env) => { const { name, uuid } = req.params; const contentRange = req.headers.get("Content-Range"); - const [start, end] = contentRange?.split("-") ?? [undefined, undefined]; + const rangeMatch = contentRange?.match(/(?:bytes\s+)?(\d+)-(\d+)/); + const [start, end] = rangeMatch ? [rangeMatch[1], rangeMatch[2]] : [undefined, undefined]; if (req.body == null) { return new Response(null, { status: 400 }); } - let contentLengthString = req.headers.get("Content-Length"); + let streamSize = getStreamSize(req.headers); let stream = req.body; - if (!contentLengthString) { + if (streamSize === undefined) { + // Without Content-Length or Content-Range the length is only known once the body has been read const blob = await req.blob(); - contentLengthString = `${blob.size}`; + streamSize = blob.size; stream = blob.stream(); } @@ -528,7 +617,7 @@ v2Router.patch("/:name+/blobs/uploads/:uuid", async (req, env: Env) => { uuid, url.pathname + "?" + url.searchParams.toString(), stream, - +contentLengthString, + streamSize, end !== undefined && start !== undefined ? [+start, +end] : undefined, ), ); @@ -546,8 +635,7 @@ v2Router.patch("/:name+/blobs/uploads/:uuid", async (req, env: Env) => { status: 202, headers: { "Location": res.location, - // Note that the HTTP Range header byte ranges are inclusive and that will be honored, even in non-standard use cases. - "Range": `${res.range.join("-")}`, + "Range": `0-${res.range[1]}`, // Ensure correct Range format (0-N) "Docker-Upload-UUID": res.id, }, }); @@ -558,15 +646,36 @@ v2Router.put("/:name+/blobs/uploads/:uuid", async (req, env: Env) => { const { digest } = req.query; const url = new URL(req.url); + let location = url.pathname + "?" + url.searchParams.toString(); + const contentLength = +(req.headers.get("Content-Length") ?? "0"); + + // A finalizing PUT may carry the last chunk. Append it through the same path a PATCH uses, so + // small chunks are combined into a valid part and an out-of-order chunk is rejected with 416 + // (instead of corrupting the assembled blob). finishUpload then completes the staged parts. + if (req.body && contentLength > 0) { + const contentRange = req.headers.get("Content-Range"); + const [start, end] = contentRange?.split("-") ?? [undefined, undefined]; + const [chunk, chunkErr] = await wrap( + env.REGISTRY_CLIENT.uploadChunk( + name, + uuid, + location, + req.body, + contentLength, + end !== undefined && start !== undefined ? [+start, +end] : undefined, + ), + ); + if (chunkErr) { + return new InternalError(); + } + if ("response" in chunk) { + return chunk.response; + } + location = chunk.location; + } + const [res, err] = await wrap( - env.REGISTRY_CLIENT.finishUpload( - name, - uuid, - url.pathname + "?" + url.searchParams.toString(), - digest! as string, - req.body ?? undefined, - +(req.headers.get("Content-Length") ?? "0"), - ), + env.REGISTRY_CLIENT.finishUpload(name, uuid, location, digest! as string), ); if (err) { @@ -590,9 +699,12 @@ v2Router.put("/:name+/blobs/uploads/:uuid", async (req, env: Env) => { v2Router.head("/:name+/blobs/:tag", async (req, env: Env) => { const { name, tag } = req.params; - const res = await env.REGISTRY.head(`${name}/blobs/${tag}`); let layerExistsResponse: CheckLayerResponse | null = null; - if (!res) { + const localResponse = await env.REGISTRY_CLIENT.layerExists(name, tag); + if ("response" in localResponse) { + return localResponse.response; + } + if (!localResponse.exists) { const registryList = registries(env); for (const registry of registryList) { const client = new RegistryHTTPClient(env, registry); @@ -610,21 +722,15 @@ v2Router.head("/:name+/blobs/:tag", async (req, env: Env) => { if (layerExistsResponse === null || !layerExistsResponse.exists) return new Response(JSON.stringify(BlobUnknownError), { status: 404 }); } else { - if (res.checksums.sha256 === null) { - throw new ServerError("invalid checksum from R2 backend"); - } - - layerExistsResponse = { - digest: hexToDigest(res.checksums.sha256!), - size: res.size, - exists: true, - }; + layerExistsResponse = localResponse; } return new Response(null, { headers: { "Content-Length": layerExistsResponse.size.toString(), "Docker-Content-Digest": layerExistsResponse.digest, + ...identityEncoding, + "Accept-Ranges": "bytes", }, }); }); @@ -686,6 +792,12 @@ v2Router.get("/:name+/tags/list", async (req, env: Env) => { v2Router.delete("/:name+/blobs/:digest", async (req, env: Env) => { const { name, digest } = req.params; + if (deletionDisabled(env)) { + return new DeletionDisabledError(); + } + if (resolveImmutableTagPattern(env.IMMUTABLE_TAG_PATTERN) !== null) { + return new ImmutableBlobError(digest); + } const res = await env.REGISTRY.head(`${name}/blobs/${digest}`); @@ -704,6 +816,9 @@ v2Router.delete("/:name+/blobs/:digest", async (req, env: Env) => { v2Router.post("/:name+/gc", async (req, env: Env) => { const { name } = req.params; + if (deletionDisabled(env)) { + return new DeletionDisabledError(); + } const mode = req.query.mode ?? "unreferenced"; if (mode !== "unreferenced" && mode !== "untagged") { diff --git a/src/utils.ts b/src/utils.ts index 1f53c6b..a045101 100644 --- a/src/utils.ts +++ b/src/utils.ts @@ -75,3 +75,25 @@ export function base64UrlDecode(s: string): string { export function base64UrlEncode(s: string): string { return Buffer.from(s, "utf8").toString("base64url"); } + +/** + * Get the estimated size of the stream (if possible). + * Does not wait for the entire stream, only checks if known length information is available. + */ +export function getStreamSize(headers: Headers): number | undefined { + const contentLength = headers.get("Content-Length"); + if (contentLength) { + return +contentLength; + } + + const contentRange = headers.get("Content-Range"); + if (contentRange) { + // Supported formats: 'bytes 0-123/456', 'bytes 0-123/*', '0-123' + const match = contentRange.match(/(?:bytes\s+)?(\d+)-(\d+)/); + if (match) { + return parseInt(match[2], 10) - parseInt(match[1], 10) + 1; + } + } + + return undefined; +} diff --git a/src/v2-errors.ts b/src/v2-errors.ts index 1fd960f..bff17d8 100644 --- a/src/v2-errors.ts +++ b/src/v2-errors.ts @@ -22,3 +22,16 @@ export const BlobUnknownError = { }, ], }; + +export const DigestInvalidError = (message = "provided digest did not match uploaded content") => + ({ + errors: [ + { + code: "DIGEST_INVALID", + message, + detail: { + message: "The provided digest did not match the content received by the registry.", + }, + }, + ], + }) as const; diff --git a/test/index.test.ts b/test/index.test.ts index b191639..dc172f7 100644 --- a/test/index.test.ts +++ b/test/index.test.ts @@ -13,6 +13,7 @@ import worker from "../index"; import { env } from "cloudflare:workers"; import { createExecutionContext, reset, waitOnExecutionContext } from "cloudflare:test"; import { base64UrlEncode } from "../src/utils"; +import { anonymousPullPatterns, anonymousPullRepository } from "../src/anonymous"; afterEach(async () => { await reset(); @@ -311,6 +312,47 @@ async function seedReferrerIndex(name: string, subjectDigest: string, descriptor } describe("v2 manifests", () => { + test("PUT /v2/:name/manifests/:reference infers mediaType from Content-Type", async () => { + const name = "helm-chart"; + const bindings = env as Env; + // Helm omits the OPTIONAL top-level mediaType; the Content-Type header carries it instead. + const withoutMediaType: Record = { ...getImageManifestV2(await generateManifest(name)) }; + delete withoutMediaType.mediaType; + + const data = JSON.stringify(withoutMediaType); + const sha256 = await getSHA256(data); + const response = await fetch( + createRequest("PUT", `/v2/${name}/manifests/v1`, new Blob([data]).stream(), { + "Content-Type": "application/vnd.oci.image.manifest.v1+json", + }), + ); + + expect(response.ok).toBeTruthy(); + expect(response.headers.get("docker-content-digest")).toEqual(sha256); + + // The stored bytes must be exactly what was pushed, or they no longer hash to the digest. + const stored = await bindings.REGISTRY.get(`${name}/manifests/${sha256}`); + expect(await stored?.text()).toEqual(data); + }); + + test("PUT /v2/:name/manifests/:reference rejects when no mediaType is available at all", async () => { + const name = "helm-chart-unknown-content-type"; + const withoutMediaType: Record = { ...getImageManifestV2(await generateManifest(name)) }; + delete withoutMediaType.mediaType; + + const response = await fetch( + createRequest("PUT", `/v2/${name}/manifests/v1`, new Blob([JSON.stringify(withoutMediaType)]).stream(), { + "Content-Type": "application/gzip", + }), + ); + + expect(response.status).toEqual(400); + const body = (await response.json()) as { errors: { code: string; message: string }[] }; + expect(body.errors[0].code).toEqual("MANIFEST_INVALID"); + // the union error must name the offending field rather than a bare "Invalid input" + expect(body.errors[0].message).toContain("mediaType"); + }); + test("HEAD /v2/:name/manifests/:reference NOT FOUND", async () => { const response = await fetch(createRequest("GET", "/v2/notfound/manifests/reference", null)); expect(response.status).toBe(404); @@ -344,10 +386,312 @@ describe("v2 manifests", () => { "content-length": "2", "content-type": "application/gzip", "docker-content-digest": sha256, + "content-encoding": "identity", }); await bindings.REGISTRY.delete(`${name}/manifests/${reference}`); }); + test("immutable release tag rejects a different manifest without changing its digest", async () => { + const bindings = env as Env & { IMMUTABLE_TAG_PATTERN?: string }; + const previousPattern = bindings.IMMUTABLE_TAG_PATTERN; + bindings.IMMUTABLE_TAG_PATTERN = String.raw`^v(?:0|[1-9]\d*)\.(?:0|[1-9]\d*)\.(?:0|[1-9]\d*)$`; + const name = "immutable-release-overwrite"; + const tag = "v1.2.3"; + + try { + const firstManifest = await generateManifest(name); + const { sha256: firstDigest } = await createManifest(name, firstManifest, tag); + const secondManifest = await generateManifest(name); + const response = await fetch( + createRequest("PUT", `/v2/${name}/manifests/${tag}`, new Blob([JSON.stringify(secondManifest)]).stream(), { + "Content-Type": "application/gzip", + }), + ); + + expect(response.status).toBe(409); + expect(await response.json()).toEqual({ + errors: [expect.objectContaining({ code: "DENIED" })], + }); + + const stored = await fetch(createRequest("HEAD", `/v2/${name}/manifests/${tag}`, null)); + expect(stored.headers.get("docker-content-digest")).toBe(firstDigest); + } finally { + bindings.IMMUTABLE_TAG_PATTERN = previousPattern; + } + }); + + test("concurrent writers cannot assign different manifests to one immutable release tag", async () => { + const bindings = env as Env & { IMMUTABLE_TAG_PATTERN?: string }; + const previousPattern = bindings.IMMUTABLE_TAG_PATTERN; + bindings.IMMUTABLE_TAG_PATTERN = String.raw`^v(?:0|[1-9]\d*)\.(?:0|[1-9]\d*)\.(?:0|[1-9]\d*)$`; + const name = "immutable-release-concurrent"; + const tag = "v1.2.3"; + + try { + const manifests = [await generateManifest(name), await generateManifest(name)]; + const manifestData = manifests.map((manifest) => JSON.stringify(manifest)); + const expectedDigests = await Promise.all(manifestData.map((data) => getSHA256(data))); + const responses = await Promise.all( + manifestData.map((data) => + fetch( + createRequest("PUT", `/v2/${name}/manifests/${tag}`, new Blob([data]).stream(), { + "Content-Type": "application/gzip", + }), + ), + ), + ); + + expect(responses.map((response) => response.status).sort()).toEqual([201, 409]); + const stored = await fetch(createRequest("HEAD", `/v2/${name}/manifests/${tag}`, null)); + expect(expectedDigests).toContain(stored.headers.get("docker-content-digest")); + } finally { + bindings.IMMUTABLE_TAG_PATTERN = previousPattern; + } + }); + + test("immutable release tag accepts an idempotent retry of the same manifest", async () => { + const bindings = env as Env & { IMMUTABLE_TAG_PATTERN?: string }; + const previousPattern = bindings.IMMUTABLE_TAG_PATTERN; + bindings.IMMUTABLE_TAG_PATTERN = String.raw`^v(?:0|[1-9]\d*)\.(?:0|[1-9]\d*)\.(?:0|[1-9]\d*)$`; + const name = "immutable-release-retry"; + const tag = "v1.2.3"; + + try { + const manifest = await generateManifest(name); + const first = await uploadManifest(name, manifest, tag); + const retry = await fetch( + createRequest("PUT", `/v2/${name}/manifests/${tag}`, new Blob([JSON.stringify(manifest)]).stream(), { + "Content-Type": "application/gzip", + }), + ); + + expect(retry.status).toBe(201); + expect(retry.headers.get("docker-content-digest")).toBe(first.sha256); + } finally { + bindings.IMMUTABLE_TAG_PATTERN = previousPattern; + } + }); + + test("tag outside the immutable pattern remains mutable", async () => { + const bindings = env as Env & { IMMUTABLE_TAG_PATTERN?: string }; + const previousPattern = bindings.IMMUTABLE_TAG_PATTERN; + bindings.IMMUTABLE_TAG_PATTERN = String.raw`^v(?:0|[1-9]\d*)\.(?:0|[1-9]\d*)\.(?:0|[1-9]\d*)$`; + const name = "mutable-operational-tag"; + const tag = "latest"; + + try { + const firstManifest = await generateManifest(name); + const { sha256: firstDigest } = await createManifest(name, firstManifest, tag); + const secondManifest = await generateManifest(name); + const { sha256: secondDigest } = await createManifest(name, secondManifest, tag); + + expect(secondDigest).not.toBe(firstDigest); + const stored = await fetch(createRequest("HEAD", `/v2/${name}/manifests/${tag}`, null)); + expect(stored.headers.get("docker-content-digest")).toBe(secondDigest); + } finally { + bindings.IMMUTABLE_TAG_PATTERN = previousPattern; + } + }); + + test("immutable release tag cannot be deleted", async () => { + const bindings = env as Env & { IMMUTABLE_TAG_PATTERN?: string }; + const previousPattern = bindings.IMMUTABLE_TAG_PATTERN; + bindings.IMMUTABLE_TAG_PATTERN = String.raw`^v(?:0|[1-9]\d*)\.(?:0|[1-9]\d*)\.(?:0|[1-9]\d*)$`; + const name = "immutable-release-delete"; + const tag = "v1.2.3"; + + try { + const manifest = await generateManifest(name); + const { sha256 } = await createManifest(name, manifest, tag); + const response = await fetch(createRequest("DELETE", `/v2/${name}/manifests/${tag}`, null)); + + expect(response.status).toBe(409); + expect(await response.json()).toEqual({ + errors: [expect.objectContaining({ code: "DENIED" })], + }); + const stored = await fetch(createRequest("HEAD", `/v2/${name}/manifests/${tag}`, null)); + expect(stored.headers.get("docker-content-digest")).toBe(sha256); + } finally { + bindings.IMMUTABLE_TAG_PATTERN = previousPattern; + } + }); + + test("tag outside the immutable pattern remains deletable", async () => { + const bindings = env as Env & { IMMUTABLE_TAG_PATTERN?: string }; + const previousPattern = bindings.IMMUTABLE_TAG_PATTERN; + bindings.IMMUTABLE_TAG_PATTERN = String.raw`^v(?:0|[1-9]\d*)\.(?:0|[1-9]\d*)\.(?:0|[1-9]\d*)$`; + const name = "mutable-operational-tag-delete"; + const tag = "latest"; + + try { + await createManifest(name, await generateManifest(name), tag); + const response = await fetch(createRequest("DELETE", `/v2/${name}/manifests/${tag}`, null)); + + expect(response.status).toBe(202); + expect((await fetch(createRequest("HEAD", `/v2/${name}/manifests/${tag}`, null))).status).toBe(404); + } finally { + bindings.IMMUTABLE_TAG_PATTERN = previousPattern; + } + }); + + test("manifest digest cannot be deleted while an immutable tag points to it", async () => { + const bindings = env as Env & { IMMUTABLE_TAG_PATTERN?: string }; + const previousPattern = bindings.IMMUTABLE_TAG_PATTERN; + bindings.IMMUTABLE_TAG_PATTERN = String.raw`^v(?:0|[1-9]\d*)\.(?:0|[1-9]\d*)\.(?:0|[1-9]\d*)$`; + const name = "immutable-release-digest-delete"; + const tag = "v1.2.3"; + + try { + const manifest = await generateManifest(name); + const { sha256 } = await createManifest(name, manifest, tag); + const response = await fetch(createRequest("DELETE", `/v2/${name}/manifests/${sha256}`, null)); + + expect(response.status).toBe(409); + expect(await response.json()).toEqual({ + errors: [expect.objectContaining({ code: "DENIED" })], + }); + expect((await fetch(createRequest("HEAD", `/v2/${name}/manifests/${tag}`, null))).status).toBe(200); + expect((await fetch(createRequest("HEAD", `/v2/${name}/manifests/${sha256}`, null))).status).toBe(200); + } finally { + bindings.IMMUTABLE_TAG_PATTERN = previousPattern; + } + }); + + test("digest deletion is blocked before paginated alias mutation when immutable policy is enabled", async () => { + const bindings = env as Env & { IMMUTABLE_TAG_PATTERN?: string }; + const previousPattern = bindings.IMMUTABLE_TAG_PATTERN; + bindings.IMMUTABLE_TAG_PATTERN = String.raw`^v(?:0|[1-9]\d*)\.(?:0|[1-9]\d*)\.(?:0|[1-9]\d*)$`; + const name = "immutable-release-paginated-digest-delete"; + const releaseTag = "v1.2.3"; + const operationalTag = "latest"; + + try { + const manifest = await generateManifest(name); + const { sha256 } = await createManifest(name, manifest, operationalTag); + await uploadManifest(name, manifest, releaseTag); + + const response = await fetch(createRequest("DELETE", `/v2/${name}/manifests/${sha256}?limit=1`, null)); + + expect(response.status).toBe(409); + expect((await fetch(createRequest("HEAD", `/v2/${name}/manifests/${operationalTag}`, null))).status).toBe(200); + expect((await fetch(createRequest("HEAD", `/v2/${name}/manifests/${releaseTag}`, null))).status).toBe(200); + expect((await fetch(createRequest("HEAD", `/v2/${name}/manifests/${sha256}`, null))).status).toBe(200); + } finally { + bindings.IMMUTABLE_TAG_PATTERN = previousPattern; + } + }); + + test("blob deletion cannot make an immutable release unpullable", async () => { + const bindings = env as Env & { IMMUTABLE_TAG_PATTERN?: string }; + const previousPattern = bindings.IMMUTABLE_TAG_PATTERN; + bindings.IMMUTABLE_TAG_PATTERN = String.raw`^v(?:0|[1-9]\d*)\.(?:0|[1-9]\d*)\.(?:0|[1-9]\d*)$`; + const name = "immutable-release-blob-delete"; + const tag = "v1.2.3"; + + try { + const manifest = getImageManifestV2(await generateManifest(name)); + await createManifest(name, manifest, tag); + const layerDigest = manifest.layers[0].digest; + + const response = await fetch(createRequest("DELETE", `/v2/${name}/blobs/${layerDigest}`, null)); + + expect(response.status).toBe(409); + expect(await response.json()).toEqual({ + errors: [expect.objectContaining({ code: "DENIED" })], + }); + expect((await fetch(createRequest("HEAD", `/v2/${name}/blobs/${layerDigest}`, null))).status).toBe(200); + const manifestResponse = await fetch(createRequest("GET", `/v2/${name}/manifests/${tag}`, null)); + expect(manifestResponse.status).toBe(200); + await manifestResponse.arrayBuffer(); + } finally { + bindings.IMMUTABLE_TAG_PATTERN = previousPattern; + } + }); + + test("blob deletion remains available when immutable tag policy is disabled", async () => { + const bindings = env as Env & { IMMUTABLE_TAG_PATTERN?: string }; + const previousPattern = bindings.IMMUTABLE_TAG_PATTERN; + bindings.IMMUTABLE_TAG_PATTERN = undefined; + const name = "mutable-blob-delete"; + + try { + const manifest = getImageManifestV2(await generateManifest(name)); + const layerDigest = manifest.layers[0].digest; + + const response = await fetch(createRequest("DELETE", `/v2/${name}/blobs/${layerDigest}`, null)); + + expect(response.status).toBe(202); + expect((await fetch(createRequest("HEAD", `/v2/${name}/blobs/${layerDigest}`, null))).status).toBe(404); + } finally { + bindings.IMMUTABLE_TAG_PATTERN = previousPattern; + } + }); + + test("invalid immutable tag policy fails before any manifest object is stored", async () => { + const bindings = env as Env & { IMMUTABLE_TAG_PATTERN?: string }; + const previousPattern = bindings.IMMUTABLE_TAG_PATTERN; + bindings.IMMUTABLE_TAG_PATTERN = "["; + const name = "invalid-immutable-policy"; + const tag = "v1.2.3"; + + try { + const manifest = await generateManifest(name); + const manifestData = JSON.stringify(manifest); + const digest = await getSHA256(manifestData); + const response = await fetch( + createRequest("PUT", `/v2/${name}/manifests/${tag}`, new Blob([manifestData]).stream(), { + "Content-Type": "application/gzip", + }), + ); + + expect(response.status).toBe(500); + expect(await bindings.REGISTRY.head(`${name}/manifests/${tag}`)).toBeNull(); + expect(await bindings.REGISTRY.head(`${name}/manifests/${digest}`)).toBeNull(); + } finally { + bindings.IMMUTABLE_TAG_PATTERN = previousPattern; + } + }); + + test("manifest PUT by digest rejects bytes whose computed digest does not match the URL", async () => { + const bindings = env as Env; + const name = "manifest-digest-mismatch"; + const manifest = await generateManifest(name); + const manifestData = JSON.stringify(manifest); + const computedDigest = await getSHA256(manifestData); + const differentManifest = await generateManifest(name); + const requestedDigest = await getSHA256(JSON.stringify(differentManifest)); + + const response = await fetch( + createRequest("PUT", `/v2/${name}/manifests/${requestedDigest}`, new Blob([manifestData]).stream(), { + "Content-Type": "application/gzip", + }), + ); + + expect(response.status).toBe(400); + expect(await response.json()).toEqual({ + errors: [expect.objectContaining({ code: "DIGEST_INVALID" })], + }); + expect(await bindings.REGISTRY.head(`${name}/manifests/${requestedDigest}`)).toBeNull(); + expect(await bindings.REGISTRY.head(`${name}/manifests/${computedDigest}`)).toBeNull(); + }); + + test("manifest PUT by its computed digest remains valid", async () => { + const name = "manifest-digest-match"; + const manifest = await generateManifest(name); + const manifestData = JSON.stringify(manifest); + const digest = await getSHA256(manifestData); + + const response = await fetch( + createRequest("PUT", `/v2/${name}/manifests/${digest}`, new Blob([manifestData]).stream(), { + "Content-Type": "application/gzip", + }), + ); + + expect(response.status).toBe(201); + expect(response.headers.get("docker-content-digest")).toBe(digest); + expect((await fetch(createRequest("HEAD", `/v2/${name}/manifests/${digest}`, null))).status).toBe(200); + }); + test("PUT then DELETE /v2/:name/manifests/:reference works", async () => { const { sha256 } = await createManifest("hello-world", await generateManifest("hello-world"), "hello"); const bindings = env as Env; @@ -474,12 +818,112 @@ describe("v2 manifests", () => { expect(layerC.ok).toBeTruthy(); expect(await layerB.bytes()).toEqual(sourceData); expect(await layerC.bytes()).toEqual(sourceData); + + // Check layer HEAD returns source metadata for mounted symlinks. + const layerHeadB = await fetch(createRequest("HEAD", `/v2/${repoB}/blobs/${layer}`, null)); + expect(layerHeadB.ok).toBeTruthy(); + expect(layerHeadB.headers.get("Docker-Content-Digest")).toEqual(layer); + expect(+(layerHeadB.headers.get("Content-Length") ?? "-1")).toEqual(sourceData.byteLength); + + const layerHeadC = await fetch(createRequest("HEAD", `/v2/${repoC}/blobs/${layer}`, null)); + expect(layerHeadC.ok).toBeTruthy(); + expect(layerHeadC.headers.get("Docker-Content-Digest")).toEqual(layer); + expect(+(layerHeadC.headers.get("Content-Length") ?? "-1")).toEqual(sourceData.byteLength); } } }); }); describe("v2 referrers", () => { + test("immutable tag conflict does not index the rejected referrer", async () => { + const bindings = env as Env; + const previousPattern = bindings.IMMUTABLE_TAG_PATTERN; + bindings.IMMUTABLE_TAG_PATTERN = String.raw`^v(?:0|[1-9]\d*)\.(?:0|[1-9]\d*)\.(?:0|[1-9]\d*)$`; + const name = "immutable-referrer-conflict"; + const tag = "v1.2.3"; + + try { + const subjectManifest = getImageManifestV2(await generateManifest(name)); + const { sha256: subjectDigest } = await createManifest(name, subjectManifest, "latest"); + const subject = { + mediaType: subjectManifest.mediaType, + digest: subjectDigest, + size: manifestSize(subjectManifest), + }; + const acceptedArtifact = { + ...getImageManifestV2(await generateManifest(name)), + subject, + annotations: { "org.opencontainers.image.title": "accepted" }, + } satisfies ManifestSchema; + const rejectedArtifact = { + ...getImageManifestV2(await generateManifest(name)), + subject, + annotations: { "org.opencontainers.image.title": "rejected" }, + } satisfies ManifestSchema; + const { sha256: acceptedDigest } = await createManifest(name, acceptedArtifact, tag); + const rejectedData = JSON.stringify(rejectedArtifact); + const rejectedDigest = await getSHA256(rejectedData); + + const response = await fetch( + createRequest("PUT", `/v2/${name}/manifests/${tag}`, new Blob([rejectedData]).stream(), { + "Content-Type": "application/gzip", + }), + ); + + expect(response.status).toBe(409); + const referrers = await getReferrersIndex(name, subjectDigest); + expect(referrers.body.manifests.map((descriptor) => descriptor.digest)).toEqual([acceptedDigest]); + expect(await bindings.REGISTRY.head(`${name}/_referrers/${subjectDigest}/${rejectedDigest}`)).toBeNull(); + } finally { + bindings.IMMUTABLE_TAG_PATTERN = previousPattern; + } + }); + + test("PUT with subject and an inferred mediaType still indexes referrers", async () => { + const name = "referrers-inferred-mediatype"; + const bindings = env as Env; + const subjectManifest = getImageManifestV2(await generateManifest(name)); + const { sha256: subjectDigest } = await createManifest(name, subjectManifest, "latest"); + + const artifactManifest = { + ...getImageManifestV2(await generateManifest(name)), + subject: { + mediaType: subjectManifest.mediaType, + digest: subjectDigest, + size: manifestSize(subjectManifest), + }, + } satisfies ManifestSchema; + const withoutMediaType: Record = { ...artifactManifest }; + delete withoutMediaType.mediaType; + + const data = JSON.stringify(withoutMediaType); + const artifactDigest = await getSHA256(data); + const response = await fetch( + createRequest("PUT", `/v2/${name}/manifests/${artifactDigest}`, new Blob([data]).stream(), { + "Content-Type": "application/vnd.oci.image.manifest.v1+json", + }), + ); + + expect(response.ok).toBeTruthy(); + expect(response.headers.get("oci-subject")).toEqual(subjectDigest); + + // Without substituting the inferred mediaType into the parsed manifest, this descriptor is + // written with mediaType undefined and every later read of it is silently dropped. + const expectedDescriptor = { + mediaType: "application/vnd.oci.image.manifest.v1+json", + digest: artifactDigest, + size: new Blob([data]).size, + // no explicit artifactType, so it falls back to the config mediaType + artifactType: artifactManifest.config.mediaType, + }; + const descriptorObject = await bindings.REGISTRY.get(`${name}/_referrers/${subjectDigest}/${artifactDigest}`); + expect(descriptorObject).not.toBeNull(); + expect(await descriptorObject?.json()).toEqual(expectedDescriptor); + + const referrers = await getReferrersIndex(name, subjectDigest); + expect(referrers.body.manifests).toEqual([expectedDescriptor]); + }); + test("PUT with subject indexes referrers and paginates results", async () => { const name = "referrers-index"; const bindings = env as Env; @@ -1102,7 +1546,7 @@ describe("v2 referrers", () => { expect(response.status).toEqual(400); }); - test("PUT /v2/:name/manifests/:reference rejects missing local subjects", async () => { + test("PUT /v2/:name/manifests/:reference accepts referrers pushed before their subject", async () => { const name = "referrers-missing-subject"; const bindings = env as Env; const missingSubjectDigest = numberedDigest(4500); @@ -1123,13 +1567,41 @@ describe("v2 referrers", () => { }), ); - expect(response.status).toEqual(400); - expect((await response.json()) as { errors: { code: string; message: string }[] }).toEqual({ - errors: [expect.objectContaining({ code: "BLOB_UNKNOWN", message: `unknown subject ${missingSubjectDigest}` })], - }); - expect(await bindings.REGISTRY.head(`${name}/manifests/artifact`)).toBeNull(); - expect(await bindings.REGISTRY.head(`${name}/manifests/${artifactDigest}`)).toBeNull(); - expect(await bindings.REGISTRY.head(`${name}/_referrers/${missingSubjectDigest}/${artifactDigest}`)).toBeNull(); + expect(response.status).toEqual(201); + expect(response.headers.get("OCI-Subject")).toEqual(missingSubjectDigest); + expect(await bindings.REGISTRY.head(`${name}/manifests/artifact`)).not.toBeNull(); + expect(await bindings.REGISTRY.head(`${name}/manifests/${artifactDigest}`)).not.toBeNull(); + expect(await bindings.REGISTRY.head(`${name}/_referrers/${missingSubjectDigest}/${artifactDigest}`)).not.toBeNull(); + + const referrers = await getReferrersIndex(name, missingSubjectDigest); + expect(referrers.body.manifests.map((m) => m.digest)).toEqual([artifactDigest]); + }); + + test("a referrer pushed before its subject remains discoverable after the subject arrives", async () => { + const name = "referrers-before-subject"; + const subjectManifest = getImageManifestV2(await generateManifest(name)); + const subjectData = JSON.stringify(subjectManifest); + const subjectDigest = await getSHA256(subjectData); + const artifactManifest = { + ...getImageManifestV2(await generateManifest(name)), + artifactType: "application/vnd.cloudchamber.btrfs-chain.v1", + subject: { + mediaType: subjectManifest.mediaType, + digest: subjectDigest, + size: manifestSize(subjectManifest), + }, + } satisfies ManifestSchema; + const { sha256: artifactDigest } = await createManifest(name, artifactManifest, "artifact"); + + const subjectResponse = await fetch( + createRequest("PUT", `/v2/${name}/manifests/${subjectDigest}`, new Blob([subjectData]).stream(), { + "Content-Type": subjectManifest.mediaType, + }), + ); + + expect(subjectResponse.status).toEqual(201); + const referrers = await getReferrersIndex(name, subjectDigest); + expect(referrers.body.manifests.map((descriptor) => descriptor.digest)).toContain(artifactDigest); }); test("PUT /v2/:name/manifests/:reference rejects invalid subject-bearing OCI indexes", async () => { @@ -1710,6 +2182,49 @@ describe("http client", () => { ); }); + test("test get layer forwards a suffix range and surfaces the partial content", async () => { + const name = "http-client-suffix-range"; + const digest = numberedDigest(9950); + const body = "abcdefghij"; + + envBindings = { ...bindings }; + envBindings.JWT_REGISTRY_TOKENS_PUBLIC_KEY = ""; + envBindings.PASSWORD = "world"; + envBindings.USERNAME = "hello"; + envBindings.REGISTRIES_JSON = undefined; + const blobRequests: { path: string; range: string | null }[] = []; + using _fetchSpy = vi.spyOn(globalThis, "fetch").mockImplementation(async (input, init) => { + const request = new Request(input as string | URL | Request, init); + const url = new URL(request.url); + if (url.pathname === "/v2/" || url.pathname === "/v2") { + return new Response(null, { status: 200 }); + } + + blobRequests.push({ path: url.pathname, range: request.headers.get("Range") }); + return new Response(body.slice(-4), { + status: 206, + headers: { "Content-Range": `bytes 6-9/${body.length}` }, + }); + }); + + const client = new RegistryHTTPClient(envBindings, { + registry: "https://localhost", + password_env: "PASSWORD", + username, + }); + + const res = await client.getLayer(name, digest, { suffix: 4 }); + if ("response" in res) { + expect(await res.response.json()).toEqual({ status: res.response.status }); + throw new Error("expected getLayer to return partial content"); + } + + expect(blobRequests).toEqual([{ path: `/v2/${name}/blobs/${digest}`, range: "bytes=-4" }]); + expect(res.contentRange).toEqual({ start: 6, end: 9, size: body.length }); + expect(res.size).toEqual(body.length); + expect(await new Response(res.stream).text()).toEqual(body.slice(-4)); + }); + test("test list referrers selects rel next from multi-link headers", async () => { const name = "http-client-referrers-multilink"; const subjectDigest = numberedDigest(9970); @@ -2252,6 +2767,40 @@ describe("v2 manifest-list", () => { expect(await layerLinked.text()).toEqual(await layerSource.text()); } }); + + test("sparse indexes can be pushed when only some platforms are mirrored", async () => { + const name = "sparse-index"; + const amd = await generateManifest(name); + const { sha256: amdDigest } = await createManifest(name, amd); + const missing = `sha256:${"f".repeat(64)}`; + const index = { + schemaVersion: 2, + mediaType: "application/vnd.oci.image.index.v1+json", + manifests: [ + { + mediaType: "application/vnd.oci.image.manifest.v1+json", + digest: amdDigest, + size: JSON.stringify(amd).length, + platform: { os: "linux", architecture: "amd64" }, + }, + { + mediaType: "application/vnd.oci.image.manifest.v1+json", + digest: missing, + size: 123, + platform: { os: "windows", architecture: "amd64" }, + }, + ], + }; + const put = await fetch( + createRequest("PUT", `/v2/${name}/manifests/latest`, new Blob([JSON.stringify(index)]).stream(), { + "Content-Type": "application/vnd.oci.image.index.v1+json", + }), + ); + expect(put.status).toBe(201); + expect(put.headers.get("docker-content-digest")).toBe(await getSHA256(JSON.stringify(index))); + expect((await fetch(createRequest("GET", `/v2/${name}/manifests/${amdDigest}`, null))).status).toBe(200); + expect((await fetch(createRequest("GET", `/v2/${name}/manifests/${missing}`, null))).status).toBe(404); + }); }); async function runGarbageCollector(name: string, mode: "unreferenced" | "untagged" | "both"): Promise { @@ -2510,6 +3059,290 @@ describe("garbage collector", () => { }); }); +describe("anonymous pulls", () => { + async function anonFetch(method: string, path: string, repositories: string | undefined, headers = {}) { + const ctx = createExecutionContext(); + const res = (await worker.fetch( + createRequest(method, path, null, headers), + { ...env, ANONYMOUS_PULL_REPOSITORIES: repositories } as Env, + ctx, + )) as Response; + await waitOnExecutionContext(ctx); + return res; + } + + test("patterns", () => { + const e = (v: string) => ({ ANONYMOUS_PULL_REPOSITORIES: v }) as Env; + expect(anonymousPullPatterns(e("")).length).toBe(0); + expect(anonymousPullPatterns(e(" a/b, c/* d ")).map((r) => r.source)).toEqual(["^a\\/b$", "^c\\/.*$", "^d$"]); + const r = (path: string, method = "GET") => + anonymousPullRepository(e("org/*,single"), new Request(`https://registry.com${path}`, { method })); + const digest = `sha256:${"a".repeat(64)}`; + expect(r("/v2/org/app/manifests/latest")).toBe("org/app"); + expect(r("/v2/org/deep/app/manifests/latest", "HEAD")).toBe("org/deep/app"); + expect(r(`/v2/org/app/blobs/${digest}`)).toBe("org/app"); + expect(r("/v2/org/app/tags/list")).toBe("org/app"); + expect(r(`/v2/org/app/referrers/${digest}`)).toBe("org/app"); + expect(r("/v2/single/manifests/1")).toBe("single"); + expect(r("/v2/single2/manifests/1")).toBeNull(); + expect(r("/v2/other/app/manifests/latest")).toBeNull(); + expect(r("/v2/org/app/manifests/latest", "PUT")).toBeNull(); + expect(r("/v2/org/app/blobs/uploads/some-uuid")).toBeNull(); + expect(r("/v2/org/app/blobs/uploads")).toBeNull(); + expect(r("/v2/")).toBeNull(); + expect(r("/v2/_catalog")).toBeNull(); + }); + + test("allowed repositories can be pulled without credentials", async () => { + const manifest = await generateManifest("public/app"); + const { sha256 } = await createManifest("public/app", manifest, "latest"); + const layer = getLayersFromManifest(manifest)[1]; + + expect((await anonFetch("GET", "/v2/public/app/manifests/latest", "public/*")).status).toBe(200); + expect((await anonFetch("HEAD", `/v2/public/app/manifests/${sha256}`, "public/*")).status).toBe(200); + const blob = await anonFetch("GET", `/v2/public/app/blobs/${layer}`, "public/*"); + expect(blob.status).toBe(200); + expect(blob.headers.get("docker-content-digest")).toBe(layer); + expect((await anonFetch("HEAD", `/v2/public/app/blobs/${layer}`, "public/*")).status).toBe(200); + expect((await anonFetch("GET", "/v2/public/app/tags/list", "public/*")).status).toBe(200); + // unknown objects in allowed repositories are a normal 404, not an auth error + expect((await anonFetch("GET", "/v2/public/app/manifests/missing", "public/*")).status).toBe(404); + }); + + test("everything else still needs credentials", async () => { + await createManifest("private/app", await generateManifest("private/app"), "latest"); + await createManifest("public/app", await generateManifest("public/app"), "latest"); + + // feature off + expect((await anonFetch("GET", "/v2/public/app/manifests/latest", undefined)).status).toBe(401); + expect((await anonFetch("GET", "/v2/public/app/manifests/latest", "")).status).toBe(401); + // other repository + expect((await anonFetch("GET", "/v2/private/app/manifests/latest", "public/*")).status).toBe(401); + // the ping keeps its Basic challenge so docker still sends credentials for pushes + const ping = await anonFetch("GET", "/v2/", "*"); + expect(ping.status).toBe(401); + expect(ping.headers.get("WWW-Authenticate")).toContain("Basic"); + expect((await anonFetch("GET", "/v2/_catalog", "*")).status).toBe(401); + // writes + expect((await anonFetch("POST", "/v2/public/app/blobs/uploads/", "*")).status).toBe(401); + expect((await anonFetch("PUT", "/v2/public/app/manifests/latest", "*")).status).toBe(401); + expect((await anonFetch("DELETE", "/v2/public/app/manifests/latest", "*")).status).toBe(401); + expect((await anonFetch("POST", "/v2/public/app/gc", "*")).status).toBe(401); + // wrong credentials are rejected even for public repositories + expect( + ( + await anonFetch("GET", "/v2/public/app/manifests/latest", "*", { + Authorization: usernamePasswordToAuth("hello", "wrong"), + }) + ).status, + ).toBe(401); + }); + + test("anonymous requests never use the pull fallback", () => { + const withFallback = { ...env, REGISTRIES_JSON: '[{ "registry": "https://ghcr.io" }]' } as Env; + expect(registries(withFallback).length).toBe(1); + expect(registries({ ...withFallback, ANONYMOUS_REQUEST: true }).length).toBe(0); + }); + + test("requests do not modify the shared env", async () => { + const shared = { ...env, ANONYMOUS_PULL_REPOSITORIES: "*" } as Env; + const ctx = createExecutionContext(); + await worker.fetch(createRequest("GET", "/v2/public/app/tags/list", null), shared, ctx); + await waitOnExecutionContext(ctx); + expect(shared.REGISTRY_CLIENT).toBeUndefined(); + expect(shared.ANONYMOUS_REQUEST).toBeUndefined(); + }); +}); + +describe("blob range requests", () => { + async function uploadBlob(name: string, data: string): Promise { + const sha256 = await getSHA256(data); + const res = await fetch(createRequest("POST", `/v2/${name}/blobs/uploads/`, null, {})); + expect(res.ok).toBeTruthy(); + const stream = limit(new Blob([data]).stream(), data.length); + const res2 = await fetch(createRequest("PATCH", res.headers.get("location")!, stream, {})); + expect(res2.ok).toBeTruthy(); + const last = await fetch(createRequest("PUT", res2.headers.get("location")! + "&digest=" + sha256, null, {})); + expect(last.ok).toBeTruthy(); + return sha256; + } + + const data = "0123456789abcdefghijklmnopqrstuvwxyz"; + + test("open-ended Range returns 206 partial content from the offset", async () => { + const name = "range-open"; + const digest = await uploadBlob(name, data); + + const res = await fetch(createRequest("GET", `/v2/${name}/blobs/${digest}`, null, { Range: "bytes=10-" })); + expect(res.status).toEqual(206); + expect(res.headers.get("content-range")).toEqual(`bytes 10-${data.length - 1}/${data.length}`); + expect(res.headers.get("content-length")).toEqual(`${data.length - 10}`); + expect(res.headers.get("accept-ranges")).toEqual("bytes"); + expect(await res.text()).toEqual(data.slice(10)); + }); + + test("bounded Range returns 206 partial content for the requested window", async () => { + const name = "range-bounded"; + const digest = await uploadBlob(name, data); + + const res = await fetch(createRequest("GET", `/v2/${name}/blobs/${digest}`, null, { Range: "bytes=5-14" })); + expect(res.status).toEqual(206); + expect(res.headers.get("content-range")).toEqual(`bytes 5-14/${data.length}`); + expect(res.headers.get("content-length")).toEqual("10"); + expect(await res.text()).toEqual(data.slice(5, 15)); + }); + + test("no Range header keeps the existing full 200 behavior", async () => { + const name = "range-none"; + const digest = await uploadBlob(name, data); + + const res = await fetch(createRequest("GET", `/v2/${name}/blobs/${digest}`, null)); + expect(res.status).toEqual(200); + expect(res.headers.get("content-length")).toEqual(`${data.length}`); + expect(res.headers.get("content-range")).toBeNull(); + expect(await res.text()).toEqual(data); + }); + + test("out-of-bounds Range returns 416 Range Not Satisfiable", async () => { + const name = "range-oob"; + const digest = await uploadBlob(name, data); + + const res = await fetch(createRequest("GET", `/v2/${name}/blobs/${digest}`, null, { Range: "bytes=100-200" })); + expect(res.status).toEqual(416); + expect(res.headers.get("content-range")).toEqual(`bytes */${data.length}`); + }); + + test("HEAD blob response advertises Accept-Ranges", async () => { + const name = "range-head"; + const digest = await uploadBlob(name, data); + + const res = await fetch(createRequest("HEAD", `/v2/${name}/blobs/${digest}`, null)); + expect(res.ok).toBeTruthy(); + expect(res.headers.get("accept-ranges")).toEqual("bytes"); + }); + + test("suffix Range returns 206 partial content with the last bytes", async () => { + const name = "range-suffix"; + const digest = await uploadBlob(name, data); + + const res = await fetch(createRequest("GET", `/v2/${name}/blobs/${digest}`, null, { Range: "bytes=-5" })); + expect(res.status).toEqual(206); + expect(res.headers.get("content-range")).toEqual(`bytes ${data.length - 5}-${data.length - 1}/${data.length}`); + expect(res.headers.get("content-length")).toEqual("5"); + expect(await res.text()).toEqual(data.slice(-5)); + }); + + test("suffix Range longer than the blob returns the whole blob as partial content", async () => { + const name = "range-suffix-oversized"; + const digest = await uploadBlob(name, data); + + const res = await fetch(createRequest("GET", `/v2/${name}/blobs/${digest}`, null, { Range: "bytes=-1000" })); + expect(res.status).toEqual(206); + expect(res.headers.get("content-range")).toEqual(`bytes 0-${data.length - 1}/${data.length}`); + expect(res.headers.get("content-length")).toEqual(`${data.length}`); + expect(await res.text()).toEqual(data); + }); + + test("zero-length suffix Range returns 416 Range Not Satisfiable", async () => { + const name = "range-suffix-zero"; + const digest = await uploadBlob(name, data); + + const res = await fetch(createRequest("GET", `/v2/${name}/blobs/${digest}`, null, { Range: "bytes=-0" })); + expect(res.status).toEqual(416); + expect(res.headers.get("content-range")).toEqual(`bytes */${data.length}`); + }); + + test("Range header without a start or a suffix length keeps the full 200 behavior", async () => { + const name = "range-malformed"; + const digest = await uploadBlob(name, data); + + const res = await fetch(createRequest("GET", `/v2/${name}/blobs/${digest}`, null, { Range: "bytes=-" })); + expect(res.status).toEqual(200); + expect(res.headers.get("content-range")).toBeNull(); + expect(await res.text()).toEqual(data); + }); +}); + +describe("blob range requests against fallback registries", () => { + const bindings = env as Env; + const data = "0123456789abcdefghijklmnopqrstuvwxyz"; + + test("an upstream 416 is reported to the client instead of a 404", async () => { + const name = "range-fallback-unsatisfiable"; + const digest = await getSHA256(data); + const previousRegistries = bindings.REGISTRIES_JSON; + bindings.REGISTRIES_JSON = JSON.stringify([{ registry: "https://fallback.registry" }]); + const blobRequests: { path: string; range: string | null }[] = []; + using _fetchSpy = vi.spyOn(globalThis, "fetch").mockImplementation(async (input, init) => { + const request = new Request(input as string | URL | Request, init); + const url = new URL(request.url); + if (url.pathname === "/v2/" || url.pathname === "/v2") { + return new Response(null, { status: 200 }); + } + + blobRequests.push({ path: url.pathname, range: request.headers.get("Range") }); + // The upstream holds the blob, but the requested range doesn't fit it. + const response = new Response(null, { + status: 416, + headers: { "Content-Range": `bytes */${data.length}`, "Accept-Ranges": "bytes" }, + }); + // A constructed Response has an empty url, which the client would read as a redirect. + Object.defineProperty(response, "url", { value: request.url }); + return response; + }); + + try { + const res = await fetch(createRequest("GET", `/v2/${name}/blobs/${digest}`, null, { Range: "bytes=100-200" })); + expect(blobRequests).toEqual([{ path: `/v2/${name}/blobs/${digest}`, range: "bytes=100-200" }]); + expect(res.status).toEqual(416); + expect(res.headers.get("content-range")).toEqual(`bytes */${data.length}`); + } finally { + bindings.REGISTRIES_JSON = previousRegistries; + } + }); + + test("a non-416 upstream failure still falls through to the next registry", async () => { + const name = "range-fallback-continue"; + const digest = await getSHA256(data); + const previousRegistries = bindings.REGISTRIES_JSON; + bindings.REGISTRIES_JSON = JSON.stringify([ + { registry: "https://broken.registry" }, + { registry: "https://healthy.registry" }, + ]); + const blobHosts: string[] = []; + using _fetchSpy = vi.spyOn(globalThis, "fetch").mockImplementation(async (input, init) => { + const request = new Request(input as string | URL | Request, init); + const url = new URL(request.url); + if (url.pathname === "/v2/" || url.pathname === "/v2") { + return new Response(null, { status: 200 }); + } + + blobHosts.push(url.host); + if (url.host === "broken.registry") { + const failure = new Response("boom", { status: 500 }); + // A constructed Response has an empty url, which the client would read as a redirect. + Object.defineProperty(failure, "url", { value: request.url }); + return failure; + } + + return new Response(data.slice(10), { + status: 206, + headers: { "Content-Range": `bytes 10-${data.length - 1}/${data.length}` }, + }); + }); + + try { + const res = await fetch(createRequest("GET", `/v2/${name}/blobs/${digest}`, null, { Range: "bytes=10-" })); + expect(blobHosts).toEqual(["broken.registry", "healthy.registry"]); + expect(res.status).toEqual(206); + expect(res.headers.get("content-range")).toEqual(`bytes 10-${data.length - 1}/${data.length}`); + expect(await res.text()).toEqual(data.slice(10)); + } finally { + bindings.REGISTRIES_JSON = previousRegistries; + } + }); +}); + test("docker.io", () => { const t = [ ["https://docker.io", true], @@ -2528,3 +3361,386 @@ test("docker.io", () => { } } }); + +describe("blob and manifest GET/HEAD response headers", () => { + // These responses must carry a literal Content-Length (and the manifest its Content-Type), and + // set Content-Encoding: identity so the runtime does not switch to chunked transfer-encoding and + // drop Content-Length — the regression these tests guard. + test("blob GET and HEAD return Content-Length, Content-Encoding identity, and the exact bytes", async () => { + const name = "headers/blob"; + const data = "blob-bytes-for-content-length"; + const sha256 = await getSHA256(data); + const post = await fetch(createRequest("POST", `/v2/${name}/blobs/uploads/`, null, {})); + const patch = await fetch( + createRequest("PATCH", post.headers.get("location")!, limit(new Blob([data]).stream(), data.length), {}), + ); + await fetch(createRequest("PUT", patch.headers.get("location")! + "&digest=" + sha256, null, {})); + + const get = await fetch(createRequest("GET", `/v2/${name}/blobs/${sha256}`, null)); + expect(get.status).toBe(200); + expect(get.headers.get("Content-Length")).toEqual(`${data.length}`); + expect(get.headers.get("Content-Encoding")).toEqual("identity"); + expect(await get.text()).toEqual(data); + + const head = await fetch(createRequest("HEAD", `/v2/${name}/blobs/${sha256}`, null)); + expect(head.status).toBe(200); + expect(head.headers.get("Content-Length")).toEqual(`${data.length}`); + expect(head.headers.get("Content-Encoding")).toEqual("identity"); + }); + + test("manifest GET and HEAD return Content-Length, Content-Type and Content-Encoding identity", async () => { + const name = "headers/manifest"; + const manifest = await generateManifest(name); + const { sha256 } = await createManifest(name, manifest, "v1"); + const size = new Blob([JSON.stringify(manifest)]).size; + + // Manifests are stored (uploadManifest) with this content type; GET/HEAD must echo it exactly. + const get = await fetch(createRequest("GET", `/v2/${name}/manifests/v1`, null)); + expect(get.status).toBe(200); + expect(get.headers.get("Content-Length")).toEqual(`${size}`); + expect(get.headers.get("Content-Type")).toEqual("application/gzip"); + expect(get.headers.get("Content-Encoding")).toEqual("identity"); + + const head = await fetch(createRequest("HEAD", `/v2/${name}/manifests/${sha256}`, null)); + expect(head.status).toBe(200); + expect(head.headers.get("Content-Length")).toEqual(`${size}`); + expect(head.headers.get("Content-Type")).toEqual("application/gzip"); + expect(head.headers.get("Content-Encoding")).toEqual("identity"); + }); +}); + +describe("manifest PUT digest validation", () => { + test("PUT manifest by a mismatched digest is rejected 400 DIGEST_INVALID (not stored)", async () => { + const name = "manifestdigest/mismatch"; + const manifest = await generateManifest(name); + const data = JSON.stringify(manifest); + const wrongDigest = "sha256:" + "0".repeat(64); + const res = await fetch( + createRequest("PUT", `/v2/${name}/manifests/${wrongDigest}`, new Blob([data]).stream(), { + "Content-Type": "application/vnd.oci.image.manifest.v1+json", + }), + ); + expect(res.status).toBe(400); + const body = (await res.json()) as { errors: { code: string }[] }; + expect(body.errors[0].code).toBe("DIGEST_INVALID"); + // Must NOT have been stored under the wrong digest key. + const get = await fetch(createRequest("GET", `/v2/${name}/manifests/${wrongDigest}`, null)); + expect(get.status).toBe(404); + }); + + test("PUT manifest by a malformed digest reference is rejected 400 DIGEST_INVALID", async () => { + const name = "manifestdigest/malformed"; + const manifest = await generateManifest(name); + const res = await fetch( + createRequest( + "PUT", + `/v2/${name}/manifests/sha256:baddigeststring`, + new Blob([JSON.stringify(manifest)]).stream(), + { + "Content-Type": "application/vnd.oci.image.manifest.v1+json", + }, + ), + ); + expect(res.status).toBe(400); + const body = (await res.json()) as { errors: { code: string }[] }; + expect(body.errors[0].code).toBe("DIGEST_INVALID"); + }); + + test("PUT manifest by its correct digest still succeeds (201)", async () => { + const name = "manifestdigest/correct"; + const manifest = await generateManifest(name); + const data = JSON.stringify(manifest); + const digest = await getSHA256(data); + const res = await fetch( + createRequest("PUT", `/v2/${name}/manifests/${digest}`, new Blob([data]).stream(), { + "Content-Type": "application/vnd.oci.image.manifest.v1+json", + }), + ); + expect(res.status).toBe(201); + expect(res.headers.get("docker-content-digest")).toEqual(digest); + }); +}); + +// Multi-chunk uploads (a PATCH chunk followed by a chunk in the finalizing PUT) exercise the +// small-chunk reconstruction path, which the registry only performs under full push-compatibility +// mode (PUSH_COMPATIBILITY_MODE=full); these finalize tests enable that mode. +async function fetchFullCompat(r: Request): Promise { + r.headers.append("Authorization", usernamePasswordToAuth(username, "world")); + const ctx = createExecutionContext(); + const res = await worker.fetch(r, { ...env, PUSH_COMPATIBILITY_MODE: "full" } as Env, ctx); + await waitOnExecutionContext(ctx); + return res as Response; +} + +describe("blob upload finalization error handling", () => { + test("PUT finalize with a wrong digest is rejected 400 DIGEST_INVALID (not 500)", async () => { + const name = "uploaderr/baddigest"; + const data = "the-real-content"; + const post = await fetch(createRequest("POST", `/v2/${name}/blobs/uploads/`, null, {})); + const patch = await fetch( + createRequest("PATCH", post.headers.get("location")!, limit(new Blob([data]).stream(), data.length), {}), + ); + expect(patch.status).toBe(202); + const wrongDigest = "sha256:" + "0".repeat(64); + const put = await fetch(createRequest("PUT", patch.headers.get("location")! + "&digest=" + wrongDigest, null, {})); + expect(put.status).toBe(400); + // Assert the mapped body, which guards the coupling to R2's checksum-mismatch wording: if the + // runtime changes that text, putBlob would fall through to 500 and this assertion would fail. + const body = (await put.json()) as { errors: { code: string }[] }; + expect(body.errors[0].code).toBe("DIGEST_INVALID"); + }); + + test("PUT finalize with an out-of-order final chunk is rejected 416", async () => { + const name = "uploaderr/outoforder"; + const first = "0123456789"; // 10 bytes staged at 0-9 + const post = await fetch(createRequest("POST", `/v2/${name}/blobs/uploads/`, null, {})); + const patch = await fetch( + createRequest("PATCH", post.headers.get("location")!, limit(new Blob([first]).stream(), first.length), { + "Content-Range": "0-9", + }), + ); + expect(patch.status).toBe(202); + // A final chunk whose range does not continue at byte 10 is out of order. + const finalData = "abcdefghij"; + const digest = await getSHA256(first + finalData); + const put = await fetch( + createRequest( + "PUT", + patch.headers.get("location")! + "&digest=" + digest, + limit(new Blob([finalData]).stream(), finalData.length), + { "Content-Range": "50-59", "Content-Length": `${finalData.length}` }, + ), + ); + expect(put.status).toBe(416); + }); + + test("PUT finalize carrying the final chunk assembles the blob (201) and round-trips", async () => { + const name = "uploaderr/finalchunk"; + const first = "first-part-bytes"; + const finalData = "final-chunk-bytes"; + const full = first + finalData; + const digest = await getSHA256(full); + const post = await fetchFullCompat(createRequest("POST", `/v2/${name}/blobs/uploads/`, null, {})); + const patch = await fetchFullCompat( + createRequest("PATCH", post.headers.get("location")!, limit(new Blob([first]).stream(), first.length), { + "Content-Range": `0-${first.length - 1}`, + }), + ); + expect(patch.status).toBe(202); + const put = await fetchFullCompat( + createRequest( + "PUT", + patch.headers.get("location")! + "&digest=" + digest, + limit(new Blob([finalData]).stream(), finalData.length), + { "Content-Range": `${first.length}-${full.length - 1}`, "Content-Length": `${finalData.length}` }, + ), + ); + expect(put.status).toBe(201); + // The PUT-carried final chunk must actually be stored (previously it was silently dropped). + const get = await fetchFullCompat(createRequest("GET", `/v2/${name}/blobs/${digest}`, null)); + expect(get.status).toBe(200); + expect(await get.text()).toEqual(full); + }); + + test("PUT finalize of an empty (zero-byte) blob succeeds 201", async () => { + const name = "uploaderr/empty"; + const emptyDigest = await getSHA256(""); + const post = await fetch(createRequest("POST", `/v2/${name}/blobs/uploads/`, null, {})); + const put = await fetch(createRequest("PUT", post.headers.get("location")! + "&digest=" + emptyDigest, null, {})); + expect(put.status).toBe(201); + const get = await fetch(createRequest("GET", `/v2/${name}/blobs/${emptyDigest}`, null)); + expect(get.status).toBe(200); + expect(await get.text()).toEqual(""); + }); +}); + +describe("mounted blob HEAD reports source metadata", () => { + test("HEAD of a cross-repo mounted blob returns the source digest and size, not the symlink's", async () => { + const src = "mountsrc/repo"; + const dst = "mountdst/repo"; + const data = "cross-repo-mounted-layer-bytes"; + const digest = await getSHA256(data); + + // Push the blob into the source repo. + const post = await fetch(createRequest("POST", `/v2/${src}/blobs/uploads/`, null, {})); + const patch = await fetch( + createRequest("PATCH", post.headers.get("location")!, limit(new Blob([data]).stream(), data.length), {}), + ); + const put = await fetch(createRequest("PUT", patch.headers.get("location")! + "&digest=" + digest, null, {})); + expect(put.ok).toBeTruthy(); + + // Cross-repo mount into the destination repo (stored as a symlink object). + const mount = await fetch(createRequest("POST", `/v2/${dst}/blobs/uploads/?from=${src}&mount=${digest}`, null, {})); + expect(mount.status).toBe(201); + expect(mount.headers.get("docker-content-digest")).toEqual(digest); + + // HEAD the mounted blob: must report the source blob's size + digest, not the symlink object's. + const head = await fetch(createRequest("HEAD", `/v2/${dst}/blobs/${digest}`, null)); + expect(head.status).toBe(200); + expect(head.headers.get("Content-Length")).toEqual(`${data.length}`); + expect(head.headers.get("Docker-Content-Digest")).toEqual(digest); + }); +}); + +describe("content-addressed writes", () => { + // Simulates a bucket with retention rules on its content prefixes (blobs, manifests by digest and + // referrer entries): such objects can be created, but never overwritten or deleted. + function retainedBucket(bucket: R2Bucket): R2Bucket { + const retained = (key: string) => /\/(blobs\/|manifests\/sha256:|_referrers\/)/.test(key); + return new Proxy(bucket, { + get(target, prop) { + if (prop === "put") { + return async (key: string, value: ReadableStream | string | null, options?: R2PutOptions) => { + if (retained(key) && (await target.head(key)) !== null) { + throw new Error(`put: object ${key} is retained and cannot be overwritten`); + } + return target.put(key, value, options); + }; + } + if (prop === "delete") { + return async (keys: string | string[]) => { + for (const key of Array.isArray(keys) ? keys : [keys]) { + if (retained(key)) throw new Error(`delete: object ${key} is retained`); + } + return target.delete(keys); + }; + } + const value = Reflect.get(target, prop); + return typeof value === "function" ? value.bind(target) : value; + }, + }); + } + + async function fetchRetained(r: Request): Promise { + r.headers.append("Authorization", usernamePasswordToAuth(username, "world")); + const bindings = env as Env; + const ctx = createExecutionContext(); + const res = await worker.fetch(r, { ...bindings, REGISTRY: retainedBucket(bindings.REGISTRY) } as Env, ctx); + await waitOnExecutionContext(ctx); + return res as Response; + } + + async function pushBlob(f: (r: Request) => Promise, name: string, data: string, digest: string) { + const post = await f(createRequest("POST", `/v2/${name}/blobs/uploads/`, null, {})); + expect(post.status).toEqual(202); + const put = await f( + createRequest("PUT", `${post.headers.get("location")!}&digest=${digest}`, new Blob([data]).stream(), { + "Content-Length": `${data.length}`, + }), + ); + expect(put.status).toEqual(201); + } + + test("re-pushing existing content does not rewrite it", async () => { + const name = "content-addressed/repush"; + const bindings = env as Env; + const manifest = await generateManifest(name); + const { sha256 } = await createManifest(name, manifest, "v1"); + const layer = getLayersFromManifest(manifest)[1]; + const manifestBefore = (await bindings.REGISTRY.head(`${name}/manifests/${sha256}`))!; + const layerBefore = (await bindings.REGISTRY.head(`${name}/blobs/${layer}`))!; + + const layerData = await (await bindings.REGISTRY.get(`${name}/blobs/${layer}`))!.text(); + await pushBlob(fetch, name, layerData, layer); + await createManifest(name, manifest, "v2"); + + expect((await bindings.REGISTRY.head(`${name}/manifests/${sha256}`))!.version).toEqual(manifestBefore.version); + expect((await bindings.REGISTRY.head(`${name}/blobs/${layer}`))!.version).toEqual(layerBefore.version); + const tags = (await (await fetch(createRequest("GET", `/v2/${name}/tags/list`, null))).json()) as TagsList; + expect(tags.tags).toEqual(["v1", "v2"]); + }); + + test("pushes of existing content succeed when the content keys cannot be overwritten", async () => { + const name = "content-addressed/retained"; + const mounted = "content-addressed/retained-mount"; + const bindings = env as Env; + const manifest = await generateManifest(name); + const data = JSON.stringify(manifest); + const { sha256 } = await createManifest(name, manifest, "v1"); + const layer = getLayersFromManifest(manifest)[1]; + const layerData = await (await bindings.REGISTRY.get(`${name}/blobs/${layer}`))!.text(); + const artifact = { + ...getImageManifestV2(await generateManifest(name)), + artifactType: "application/vnd.example.signature.v1", + subject: { mediaType: "application/vnd.oci.image.manifest.v1+json", digest: sha256, size: data.length }, + } satisfies ManifestSchema; + const artifactData = JSON.stringify(artifact); + const artifactDigest = await getSHA256(artifactData); + await createManifest(name, artifact); + expect(await mountLayersFromManifest(name, manifest, mounted)).toBeGreaterThan(0); + + // Push everything a second time through a bucket that refuses overwrites of content keys + await pushBlob(fetchRetained, name, layerData, layer); + for (const reference of [sha256, "v1", "v2"]) { + const res = await fetchRetained( + createRequest("PUT", `/v2/${name}/manifests/${reference}`, new Blob([data]).stream(), { + "Content-Type": "application/gzip", + }), + ); + expect(res.status).toEqual(201); + expect(res.headers.get("docker-content-digest")).toEqual(sha256); + } + const referrer = await fetchRetained( + createRequest("PUT", `/v2/${name}/manifests/${artifactDigest}`, new Blob([artifactData]).stream(), { + "Content-Type": "application/gzip", + }), + ); + expect(referrer.status).toEqual(201); + for (const digest of getLayersFromManifest(manifest)) { + const mount = await fetchRetained( + createRequest("POST", `/v2/${mounted}/blobs/uploads/?from=${name}&mount=${digest}`, null, {}), + ); + expect(mount.status).toEqual(201); + } + + const layerGet = await fetch(createRequest("GET", `/v2/${mounted}/blobs/${layer}`, null)); + expect(await layerGet.text()).toEqual(layerData); + const referrers = await getReferrersIndex(name, sha256); + expect(referrers.body.manifests.map((m) => m.digest)).toEqual([artifactDigest]); + const tags = (await (await fetch(createRequest("GET", `/v2/${name}/tags/list`, null))).json()) as TagsList; + expect(tags.tags).toEqual(["v1", "v2"]); + }); +}); + +describe("DISABLE_DELETE", () => { + async function fetchNoDelete(r: Request, value = "true"): Promise { + r.headers.append("Authorization", usernamePasswordToAuth(username, "world")); + const ctx = createExecutionContext(); + const res = await worker.fetch(r, { ...env, DISABLE_DELETE: value } as Env, ctx); + await waitOnExecutionContext(ctx); + return res as Response; + } + + test("refuses deleting manifests, blobs and garbage collection", async () => { + const name = "no-delete/app"; + const bindings = env as Env; + const manifest = await generateManifest(name); + const { sha256 } = await createManifest(name, manifest, "v1"); + const layer = getLayersFromManifest(manifest)[1]; + + for (const path of [`/v2/${name}/manifests/v1`, `/v2/${name}/manifests/${sha256}`, `/v2/${name}/blobs/${layer}`]) { + const res = await fetchNoDelete(createRequest("DELETE", path, null)); + expect(res.status).toEqual(405); + expect(((await res.json()) as { errors: { code: string }[] }).errors[0].code).toEqual("UNSUPPORTED"); + } + for (const mode of ["unreferenced", "untagged"]) { + const res = await fetchNoDelete(createRequest("POST", `/v2/${name}/gc?mode=${mode}`, null)); + expect(res.status).toEqual(405); + } + + expect(await bindings.REGISTRY.head(`${name}/manifests/v1`)).not.toBeNull(); + expect(await bindings.REGISTRY.head(`${name}/manifests/${sha256}`)).not.toBeNull(); + expect(await bindings.REGISTRY.head(`${name}/blobs/${layer}`)).not.toBeNull(); + }); + + test("leaves deletion enabled for any other value", async () => { + const name = "no-delete/off"; + const { sha256 } = await createManifest(name, await generateManifest(name), "v1"); + expect((await fetchNoDelete(createRequest("DELETE", `/v2/${name}/manifests/v1`, null), "false")).status).toEqual( + 202, + ); + expect((await fetchNoDelete(createRequest("DELETE", `/v2/${name}/manifests/${sha256}`, null), "")).status).toEqual( + 202, + ); + }); +}); diff --git a/test/r2-write-order.test.ts b/test/r2-write-order.test.ts new file mode 100644 index 0000000..f19747c --- /dev/null +++ b/test/r2-write-order.test.ts @@ -0,0 +1,101 @@ +import { describe, expect, test } from "vitest"; +import type { Env } from ".."; +import { R2Registry } from "../src/registry/r2"; + +const imageManifestContentType = "application/vnd.oci.image.manifest.v1+json"; + +describe("manifest write ordering", () => { + test("starts mutable tag and digest writes concurrently", async () => { + const startedManifestWrites: string[] = []; + let writesStartedBeforeFirstCompletion: string[] = []; + let releaseWrites: () => void; + const writesMayComplete = new Promise((resolve) => { + releaseWrites = resolve; + }); + let releaseScheduled = false; + const registry = { + head: async () => null, + put: async (key: string) => { + if (key.includes("/manifests/")) { + startedManifestWrites.push(key); + if (!releaseScheduled) { + releaseScheduled = true; + queueMicrotask(() => { + writesStartedBeforeFirstCompletion = [...startedManifestWrites]; + releaseWrites(); + }); + } + await writesMayComplete; + } + return {}; + }, + } as unknown as R2Bucket; + const env = { REGISTRY: registry } as Env; + const manifest = JSON.stringify({ + schemaVersion: 2, + mediaType: imageManifestContentType, + config: { + mediaType: "application/vnd.oci.image.config.v1+json", + digest: `sha256:${"0".repeat(64)}`, + size: 0, + }, + layers: [], + }); + + const result = await new R2Registry(env).putManifestInner( + "write-order", + "latest", + new Blob([manifest]).stream(), + imageManifestContentType, + false, + ); + + expect("response" in result).toBe(false); + expect(writesStartedBeforeFirstCompletion).toEqual([ + expect.stringMatching(/\/manifests\/sha256:/), + "write-order/manifests/latest", + ]); + }); + + test("finishes the digest write before starting a protected tag write", async () => { + const events: string[] = []; + const tagKey = "protected-write-order/manifests/v1.2.3"; + const registry = { + head: async () => null, + put: async (key: string) => { + if (key.includes("/manifests/")) { + const kind = key === tagKey ? "tag" : "digest"; + events.push(`start:${kind}`); + await Promise.resolve(); + events.push(`finish:${kind}`); + } + return {}; + }, + } as unknown as R2Bucket; + const env = { + REGISTRY: registry, + IMMUTABLE_TAG_PATTERN: String.raw`v(?:0|[1-9]\d*)\.(?:0|[1-9]\d*)\.(?:0|[1-9]\d*)`, + } as Env; + const manifest = JSON.stringify({ + schemaVersion: 2, + mediaType: imageManifestContentType, + config: { + mediaType: "application/vnd.oci.image.config.v1+json", + digest: `sha256:${"0".repeat(64)}`, + size: 0, + }, + layers: [], + }); + + const result = await new R2Registry(env).putManifestInner( + "protected-write-order", + "v1.2.3", + new Blob([manifest]).stream(), + imageManifestContentType, + false, + ); + + expect("response" in result).toBe(false); + expect(events).toEqual(["start:digest", "finish:digest", "start:tag", "finish:tag"]); + }); +}); diff --git a/test/streaming.test.ts b/test/streaming.test.ts new file mode 100644 index 0000000..eb0b285 --- /dev/null +++ b/test/streaming.test.ts @@ -0,0 +1,50 @@ +import { describe, expect, test } from "vitest"; +import { getStreamSize } from "../src/utils"; +import { limit } from "../src/chunk"; + +describe("Streaming Utilities", () => { + test("getStreamSize extracts size correctly from Content-Length", () => { + const headers = new Headers({ + "Content-Length": "1024", + }); + expect(getStreamSize(headers)).toBe(1024); + }); + + test("getStreamSize extracts size correctly from Content-Range (standard)", () => { + const headers = new Headers({ + "Content-Range": "bytes 0-52428799/104857600", + }); + expect(getStreamSize(headers)).toBe(52428800); + }); + + test("getStreamSize extracts size correctly from Content-Range (without total)", () => { + const headers = new Headers({ + "Content-Range": "bytes 500-999", + }); + expect(getStreamSize(headers)).toBe(500); + }); + + test("getStreamSize returns undefined when no size headers are present", () => { + const headers = new Headers(); + expect(getStreamSize(headers)).toBeUndefined(); + }); +}); + +describe("Stream limit function", () => { + test("limit correctly truncates a stream", async () => { + const data = new Uint8Array([1, 2, 3, 4, 5, 6, 7, 8, 9, 10]); + const stream = new Blob([data]).stream(); + const limitedStream = limit(stream, 5); + + const reader = limitedStream.getReader(); + const chunks = []; + while (true) { + const { done, value } = await reader.read(); + if (done) break; + chunks.push(...value); + } + + expect(chunks).toHaveLength(5); + expect(new Uint8Array(chunks)).toEqual(new Uint8Array([1, 2, 3, 4, 5])); + }); +}); diff --git a/test/tag-policy.test.ts b/test/tag-policy.test.ts new file mode 100644 index 0000000..a28957a --- /dev/null +++ b/test/tag-policy.test.ts @@ -0,0 +1,28 @@ +import { describe, expect, test } from "vitest"; +import { isImmutableTagReference, resolveImmutableTagPattern } from "../src/registry/tag-policy"; + +describe("immutable tag policy", () => { + test("is disabled when the pattern is unset or empty", () => { + expect(resolveImmutableTagPattern(undefined)).toBeNull(); + expect(resolveImmutableTagPattern(" ")).toBeNull(); + }); + + test("matches the entire reference and excludes operational tags and prereleases", () => { + const pattern = resolveImmutableTagPattern(String.raw`v(?:0|[1-9]\d*)\.(?:0|[1-9]\d*)\.(?:0|[1-9]\d*)`); + + expect(isImmutableTagReference("v1.2.3", pattern)).toBe(true); + expect(isImmutableTagReference("prefix-v1.2.3", pattern)).toBe(false); + expect(isImmutableTagReference("v1.2.3-suffix", pattern)).toBe(false); + expect(isImmutableTagReference("v1.2.3-rc.1", pattern)).toBe(false); + expect(isImmutableTagReference("latest", pattern)).toBe(false); + }); + + test("never treats a content digest as a tag reference", () => { + const pattern = resolveImmutableTagPattern(".*"); + expect(isImmutableTagReference(`sha256:${"a".repeat(64)}`, pattern)).toBe(false); + }); + + test("rejects an invalid configured expression", () => { + expect(() => resolveImmutableTagPattern("[")).toThrow(/Invalid IMMUTABLE_TAG_PATTERN/); + }); +});