Skip to content
Draft
Show file tree
Hide file tree
Changes from all commits
Commits
File filter

Filter by extension

Filter by extension

Conversations
Failed to load comments.
Loading
Jump to
Jump to file
Failed to load files.
Loading
Diff view
Diff view
114 changes: 100 additions & 14 deletions src/registry/r2.ts
Original file line number Diff line number Diff line change
Expand Up @@ -28,12 +28,23 @@ import {
UploadId,
UploadObject,
wrapError,
BlobRangeRequest,
} from "./registry";
import { GarbageCollectionMode, GarbageCollector } from "./garbage-collector";
import { ManifestSchema, manifestSchema } from "../manifest";

export const ociImageIndexContentType = "application/vnd.oci.image.index.v1+json";

function rangeNotSatisfiableResponse(size: number): Response {
return new Response(null, {
status: 416,
headers: {
"Content-Range": `bytes */${size}`,
"Accept-Ranges": "bytes",
},
});
}

function referrersPrefix(name: string, digest: string): string {
return `${name}/_referrers/${digest}/`;
}
Expand Down Expand Up @@ -690,38 +701,113 @@ export class R2Registry implements Registry {
};
}

async getLayer(name: string, digest: string): Promise<RegistryError | GetLayerResponse> {
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<RegistryError | GetLayerResponse> {
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;
const start = range.offset;
if (start < 0 || start >= totalSize) {
return { response: rangeNotSatisfiableResponse(totalSize) };
}

const 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<RegistryError | GetLayerResponse> {
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<RegistryError | UploadObject> {
// Generate a unique ID for this upload
const uuid = crypto.randomUUID();
Expand Down
10 changes: 9 additions & 1 deletion src/registry/registry.ts
Original file line number Diff line number Diff line change
Expand Up @@ -99,11 +99,19 @@ export type GetManifestResponse = {
contentType: string;
};

// requested byte range for a layer read (inclusive end; end omitted means until the end of the object)
export type BlobRangeRequest = {
offset: number;
end?: 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 = {
Expand Down Expand Up @@ -150,7 +158,7 @@ export interface Registry {
layerExists(namespace: string, digest: string): Promise<CheckLayerResponse | RegistryError>;

// get a layer stream from the registry
getLayer(namespace: string, digest: string): Promise<GetLayerResponse | RegistryError>;
getLayer(namespace: string, digest: string, range?: BlobRangeRequest): Promise<GetLayerResponse | RegistryError>;

// list referrers for a subject digest
listReferrers(
Expand Down
54 changes: 41 additions & 13 deletions src/router.ts
Original file line number Diff line number Diff line change
Expand Up @@ -358,16 +358,48 @@ v2Router.get("/:name+/referrers/:digest", async (req, env: Env) => {
);
});

// Parses a single HTTP byte range request of the form "bytes=<start>-" or "bytes=<start>-<end>".
// Multi-range, suffix ("bytes=-<n>") and malformed values are ignored so the full object is served.
function parseBlobRange(header: string | null): { offset: number; end?: number } | undefined {
if (header === null) return undefined;
const match = /^bytes=(\d+)-(\d*)$/.exec(header.trim());
if (match === null) return undefined;
const offset = Number(match[1]);
if (!Number.isInteger(offset)) return undefined;
if (match[2] === "") return { offset };
const end = Number(match[2]);
if (!Number.isInteger(end)) return { offset };
return { offset, end };
}

function blobGetResponse(layer: GetLayerResponse): Response {
const headers: Record<string, string> = {
"Docker-Content-Digest": layer.digest,
"Accept-Ranges": "bytes",
};
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;
Expand Down Expand Up @@ -400,12 +432,7 @@ v2Router.get("/:name+/blobs/:digest", async (req, env: Env, context: ExecutionCo

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) => {
Expand Down Expand Up @@ -625,6 +652,7 @@ v2Router.head("/:name+/blobs/:tag", async (req, env: Env) => {
headers: {
"Content-Length": layerExistsResponse.size.toString(),
"Docker-Content-Digest": layerExistsResponse.digest,
"Accept-Ranges": "bytes",
},
});
});
Expand Down
68 changes: 68 additions & 0 deletions test/index.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -2362,6 +2362,74 @@ describe("garbage collector", () => {
});
});

describe("blob range requests", () => {
async function uploadBlob(name: string, data: string): Promise<string> {
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("docker.io", () => {
const t = [
["https://docker.io", true],
Expand Down