Skip to content
Merged
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
60 changes: 30 additions & 30 deletions src/core/cliManager.ts
Original file line number Diff line number Diff line change
Expand Up @@ -805,50 +805,50 @@ export class CliManager {
: (buffer.byteLength / contentLength) * 100,
});
if (onProgress) {
progressWrite = onProgress(
written,
Number.isNaN(contentLength) ? null : contentLength,
).catch((error) => {
this.output.warn(
"Failed to write progress log:",
errToStr(error),
);
});
// Chain so awaiting the final progressWrite awaits every
// progress-log write, not just the last one.
progressWrite = progressWrite
.then(() =>
onProgress(
written,
Number.isNaN(contentLength) ? null : contentLength,
),
)
.catch((error) => {
this.output.warn(
"Failed to write progress log:",
errToStr(error),
);
});
}
});
});

// Wait for the stream to end or error.
return new Promise<boolean>((resolve, reject) => {
// fs emits "close" only after every pending write callback, so
// settling after it cannot race the trailing progress-log write.
const settle = (settleFn: () => void): void => {
writeStream.once("close", () => {
void progressWrite.then(settleFn);
});
};
const downloadError = (error: unknown): Error =>
new Error(
`Unable to download binary: ${errToStr(error, "no reason given")}`,
);

writeStream.on("error", (error) => {
readStream.destroy();
void progressWrite.then(() =>
reject(
new Error(
`Unable to download binary: ${errToStr(error, "no reason given")}`,
),
),
);
settle(() => reject(downloadError(error)));
});
readStream.on("error", (error) => {
writeStream.close();
void progressWrite.then(() =>
reject(
new Error(
`Unable to download binary: ${errToStr(error, "no reason given")}`,
),
),
);
settle(() => reject(downloadError(error)));
});
readStream.on("close", () => {
writeStream.close();
void progressWrite.then(() => {
if (cancelled) {
resolve(false);
} else {
resolve(true);
}
});
settle(() => resolve(!cancelled));
});
});
},
Expand Down
10 changes: 4 additions & 6 deletions test/unit/core/cliManager.concurrent.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -93,14 +93,12 @@ function setupManager(testDir: string): CliManager {
}

/**
* Asserts the lock and progress files are removed. The lock directory can
* briefly reappear while a peer re-acquires and releases it, so poll until gone.
* Asserts the lock and progress files are removed. fetchBinary settles only
* after trailing fs writes, so no polling is needed.
*/
async function expectLockFilesRemoved(binaryPath: string): Promise<void> {
await vi.waitFor(async () => {
await expect(fs.access(binaryPath + ".lock")).rejects.toThrow();
await expect(fs.access(binaryPath + ".progress.log")).rejects.toThrow();
});
await expect(fs.access(binaryPath + ".lock")).rejects.toThrow();
await expect(fs.access(binaryPath + ".progress.log")).rejects.toThrow();
}

describe("CliManager Concurrent Downloads", () => {
Expand Down
8 changes: 8 additions & 0 deletions test/unit/core/cliManager.test.ts
Original file line number Diff line number Diff line change
Expand Up @@ -707,6 +707,14 @@ describe("CliManager", () => {
}
},
);

it("waits for trailing fs writes before finishing, keeping the progress log cleaned up", async () => {
const { manager, mockApi, withTrailingWriteFlush } = setupCliManager();
withTrailingWriteFlush();
expectPathsEqual(await manager.fetchBinary(mockApi), BINARY_PATH);
await flushPendingIO();
expect(memfs.existsSync(`${BINARY_PATH}.progress.log`)).toBe(false);
});
});

describe("Download Progress Tracking", () => {
Expand Down
33 changes: 32 additions & 1 deletion test/unit/core/cliManagerHarness.ts
Original file line number Diff line number Diff line change
Expand Up @@ -112,6 +112,32 @@ export function setupCliManager(basePath: string = BASE_PATH) {
createMockStream(partial, { error: new Error("connection reset") }),
);

/**
* Serve a download whose fs write callbacks are held until close() is
* called, like a real stream on a loaded runner.
*/
const withTrailingWriteFlush = () => {
vi.spyOn(fs, "createWriteStream").mockImplementation((writePath) => {
const pendingWrites: Array<() => void> = [];
const stream = new EventEmitter() as fs.WriteStream;
stream.write = ((chunk: Buffer, callback?: () => void) => {
memfs.appendFileSync(String(writePath), chunk);
if (callback) {
pendingWrites.push(callback);
}
return true;
}) as fs.WriteStream["write"];
stream.close = () => {
setImmediate(() => {
pendingWrites.splice(0).forEach((flush) => flush());
setImmediate(() => stream.emit("close"));
});
};
return stream;
});
withSuccessfulDownload();
};

/** Queue one HTTP response per signature source status. */
const withSignatureResponses = (statuses: number[]) => {
for (const status of statuses) {
Expand All @@ -134,7 +160,11 @@ export function setupCliManager(basePath: string = BASE_PATH) {
const stream = new EventEmitter();
(stream as unknown as fs.WriteStream).write = vi.fn();
(stream as unknown as fs.WriteStream).close = vi.fn();
setImmediate(() => stream.emit("error", new Error(message)));
setImmediate(() => {
stream.emit("error", new Error(message));
// A real fs stream with autoClose emits "close" after "error".
setImmediate(() => stream.emit("close"));
});
return stream as ReturnType<typeof memfs.createWriteStream>;
});
withHttpResponse(
Expand Down Expand Up @@ -183,6 +213,7 @@ export function setupCliManager(basePath: string = BASE_PATH) {
withSignatureResponses,
withInvalidSignature,
withStreamError,
withTrailingWriteFlush,
event,
noEvent,
expectProps,
Expand Down
Loading