Skip to content
Closed
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
12 changes: 12 additions & 0 deletions src/file-store.gcs.spec.ts
Original file line number Diff line number Diff line change
Expand Up @@ -118,10 +118,22 @@ describe("GCPFileStore", () => {
describe("getAsStream", () => {
it("should return the read stream", async () => {
const rs = Readable.from(["data"]);
mockFile.exists.mockResolvedValueOnce([true]);
mockFile.createReadStream.mockReturnValueOnce(rs);

expect(await store.getAsStream("test.txt")).toBe(rs);
});

// createReadStream does no I/O, so without an existence check a missing object
// resolves and fails later as a stream 'error'. Every other provider rejects.
it("rejects for a missing object instead of returning a doomed stream", async () => {
mockFile.exists.mockResolvedValueOnce([false]);

await expect(store.getAsStream("missing.txt")).rejects.toThrow(
"File not found: missing.txt",
);
expect(mockFile.createReadStream).not.toHaveBeenCalled();
});
});

describe("copyFromStream", () => {
Expand Down
32 changes: 22 additions & 10 deletions src/file-store.local.spec.ts
Original file line number Diff line number Diff line change
Expand Up @@ -75,6 +75,19 @@ describe("LocalFileStore", () => {
fastify.register(FileStorePlugin, { type: "nope" as any }),
).rejects.toThrow("Unknown storage type: nope");
});

// A plain lookup finds Object.prototype.toString and registers with no FileStore
// decorated — a typo becoming a silent no-op boot instead of a failure.
it.each(["toString", "constructor", "valueOf", "hasOwnProperty"])(
"rejects the inherited Object.prototype name %s",
async (type) => {
const f = Fastify({ logger: false });
await expect(
f.register(FileStorePlugin, { type: type as never }),
).rejects.toThrow(`Unknown storage type: ${type}`);
await f.close();
},
);
});

describe("exists", () => {
Expand Down Expand Up @@ -165,12 +178,11 @@ describe("LocalFileStore", () => {
expect((await store.getAsBuffer("read.txt")).toString()).toBe("buffered");
});

it("throws File not found when stat yields nothing", async () => {
it("rejects with ENOENT for a missing file", async () => {
const store = await register();
jest.spyOn(fs.promises, "stat").mockResolvedValue(undefined);
await expect(store.getAsBuffer("ghost.txt")).rejects.toThrow(
`File not found: ${path.join(tempDir, "ghost.txt")}`,
);
await expect(store.getAsBuffer("ghost.txt")).rejects.toMatchObject({
code: "ENOENT",
});
});
});

Expand All @@ -182,12 +194,12 @@ describe("LocalFileStore", () => {
expect((await streamToBuffer(rs)).toString()).toBe("streamed");
});

it("throws File not found when stat yields nothing", async () => {
// Rejects rather than deferring to a stream 'error': createReadStream is lazy.
it("rejects with ENOENT for a missing file", async () => {
const store = await register();
jest.spyOn(fs.promises, "stat").mockResolvedValue(undefined);
await expect(store.getAsStream("ghost.txt")).rejects.toThrow(
`File not found: ${path.join(tempDir, "ghost.txt")}`,
);
await expect(store.getAsStream("ghost.txt")).rejects.toMatchObject({
code: "ENOENT",
});
});
});

Expand Down
31 changes: 21 additions & 10 deletions src/file-store.ts
Original file line number Diff line number Diff line change
Expand Up @@ -79,7 +79,11 @@ const Configure = {
const plugin: FastifyPluginAsync<{
type: keyof typeof Configure;
}> = async function (f, opts): Promise<void> {
const configure = Configure[opts.type];
// hasOwn, not a plain lookup: `Configure["toString"]` finds Object.prototype.toString
// and registers with no FileStore decorated, turning a typo into a silent no-op boot.
const configure = Object.prototype.hasOwnProperty.call(Configure, opts.type)
? Configure[opts.type]
: undefined;
if (!configure) {
throw new Error(`Unknown storage type: ${opts.type}`);
}
Expand Down Expand Up @@ -127,11 +131,10 @@ class LocalFileStore implements FileStore {
await fs.promises.writeFile(p, data);
}
async getAsBuffer(filepath: string): Promise<Buffer> {
const p = path.join(this.dir, filepath);
if (await fs.promises.stat(p)) {
return fs.promises.readFile(p);
}
throw new Error(`File not found: ${p}`);
// No stat guard: a fulfilled stat is never falsy, so the old "File not found"
// branch was unreachable and readFile already rejects with ENOENT. Dropping it
// also removes a redundant stat and the TOCTOU window between the two calls.
return fs.promises.readFile(path.join(this.dir, filepath));
}
async copyFromLocalFile(
filepath: string,
Expand All @@ -145,10 +148,10 @@ class LocalFileStore implements FileStore {
}
async getAsStream(filepath: string): Promise<NodeJS.ReadableStream> {
const p = path.join(this.dir, filepath);
if (await fs.promises.stat(p)) {
return fs.createReadStream(p);
}
throw new Error(`File not found: ${p}`);
// stat first so a missing file rejects the promise rather than surfacing later as
// a stream 'error' — createReadStream is lazy, like the GCS reader.
await fs.promises.stat(p);
return fs.createReadStream(p);
}
async copyFromStream(
filepath: string,
Expand Down Expand Up @@ -302,6 +305,14 @@ class GCPFileStore implements FileStore {

async getAsStream(filepath: string): Promise<NodeJS.ReadableStream> {
const gcsfile = this.storage.bucket(this.bucket).file(filepath);
// createReadStream does no I/O before returning, so without this check a missing
// object resolves and only fails later as a stream 'error' — which a caller doing
// `reply.send(await getAsStream(k))` turns into a 200 with a broken body, or an
// unhandled 'error' that takes the process down. Every other provider rejects.
const [exists] = await gcsfile.exists();
if (!exists) {
throw new Error(`File not found: ${filepath}`);
}
return gcsfile.createReadStream();
}

Expand Down
58 changes: 50 additions & 8 deletions src/utils.spec.ts
Original file line number Diff line number Diff line change
@@ -1,3 +1,4 @@
import { EventEmitter } from "node:events";
import { Readable } from "node:stream";
import { DataStream, streamToBuffer } from "./utils";

Expand Down Expand Up @@ -163,15 +164,56 @@ describe("Utils", () => {
});
});

describe("DataStream interface", () => {
it("should define the correct interface structure", () => {
// This is a compile-time test to ensure the interface is properly defined
const mockStream: DataStream = {
on: jest.fn().mockReturnThis(),
};
describe("premature close", () => {
// A stream destroyed without an error emits only 'close'. Before this was handled
// the promise stayed pending forever and hung the request awaiting getAsBuffer.
it("rejects when the stream is destroyed without an error", async () => {
const stream = new Readable({ read() {} });
stream.push("partial");
setImmediate(() => stream.destroy());

await expect(streamToBuffer(stream)).rejects.toThrow(
"stream closed before end",
);
});

it("rejects with the original error when destroyed with one", async () => {
const stream = new Readable({ read() {} });
setImmediate(() => stream.destroy(new Error("boom")));

await expect(streamToBuffer(stream)).rejects.toThrow("boom");
});

it("still resolves normally when close follows end", async () => {
const stream: DataStream = Readable.from(["a", "b"]);

expect((await streamToBuffer(stream)).toString()).toBe("ab");
});

// Late events must not settle the promise a second time. EventEmitter is used
// directly so the exact sequence can be driven; a real Readable will not emit
// 'error' after 'end'.
it("ignores an error arriving after end", async () => {
const em = new EventEmitter() as unknown as DataStream;
const p = streamToBuffer(em);
const e = em as unknown as EventEmitter;
e.emit("data", Buffer.from("ok"));
e.emit("end");
e.emit("error", new Error("too late"));
e.emit("close");

expect((await p).toString()).toBe("ok");
});

it("ignores end and close arriving after an error", async () => {
const em = new EventEmitter() as unknown as DataStream;
const p = streamToBuffer(em);
const e = em as unknown as EventEmitter;
e.emit("error", new Error("first"));
e.emit("end");
e.emit("close");

expect(typeof mockStream.on).toBe("function");
expect(mockStream.on("test", () => {})).toBe(mockStream);
await expect(p).rejects.toThrow("first");
});
});
});
24 changes: 20 additions & 4 deletions src/utils.ts
Original file line number Diff line number Diff line change
Expand Up @@ -8,13 +8,29 @@ export interface DataStream {
export function streamToBuffer(stream: DataStream): Promise<Buffer> {
return new Promise<Buffer>((resolve, reject) => {
const chunks: Buffer[] = [];
let settled = false;

stream.on("data", (chunk: Buffer | string) =>
chunks.push(Buffer.isBuffer(chunk) ? chunk : Buffer.from(chunk)),
);
// Typed as Error so the rejection carries a stack; node streams always emit one.
stream.on("error", (err: Error) => reject(err));
stream.on("end", () =>
resolve(chunks.length === 1 ? chunks[0] : Buffer.concat(chunks)),
);
stream.on("error", (err: Error) => {
if (settled) return;
settled = true;
reject(err);
});
stream.on("end", () => {
if (settled) return;
settled = true;
resolve(chunks.length === 1 ? chunks[0] : Buffer.concat(chunks));
});
// A stream destroyed without an error argument emits only 'close', which would
// otherwise leave this promise pending forever and hang the awaiting request.
// 'close' follows 'end' in the normal path, where the settled flag ignores it.
stream.on("close", () => {
if (settled) return;
settled = true;
reject(new Error("stream closed before end"));
});
});
}
Loading