🐛 back: zip file creation throttling

This commit is contained in:
Eric Doughty-Papassideris
2025-01-31 06:23:13 +01:00
parent 988907e7c3
commit 4444e73bd1
2 changed files with 100 additions and 17 deletions
@@ -49,6 +49,7 @@ import {
isVirtualFolder, isVirtualFolder,
updateItemSize, updateItemSize,
isInTrash, isInTrash,
ThrottledArchive,
} from "../utils"; } from "../utils";
import { getConfigOrDefault } from "../../../utils/get-config"; import { getConfigOrDefault } from "../../../utils/get-config";
import { import {
@@ -699,9 +700,8 @@ export class DocumentsService {
} }
return item; return item;
} catch (error) { } catch (err) {
console.error(error); this.logger.error({ err }, "Failed to update drive item");
this.logger.error({ error: `${error}` }, "Failed to update drive item");
throw new CrudException("Failed to update item", 500); throw new CrudException("Failed to update item", 500);
} }
}; };
@@ -1567,6 +1567,10 @@ export class DocumentsService {
zlib: { level: 9 }, zlib: { level: 9 },
}); });
const throttledArchive = new ThrottledArchive(archive, async ({ readable, name, prefix }) => {
archive.append(readable, { name, prefix });
});
archive.on("error", error => { archive.on("error", error => {
this.logger.error("error while creating ZIP file: ", error); this.logger.error("error while creating ZIP file: ", error);
}); });
@@ -1582,7 +1586,7 @@ export class DocumentsService {
await addDriveItemToArchive( await addDriveItemToArchive(
id, id,
null, null,
archive, throttledArchive,
archive => { archive => {
if (didBeginTransmission) return; if (didBeginTransmission) return;
didBeginTransmission = true; didBeginTransmission = true;
@@ -1597,6 +1601,8 @@ export class DocumentsService {
} }
if (!didBeginTransmission) beginArchiveTransmit(archive); if (!didBeginTransmission) beginArchiveTransmit(archive);
await throttledArchive.waitUntilEmptied();
//TODO[ASH] why do we need this call?? //TODO[ASH] why do we need this call??
archive.finalize(); archive.finalize();
@@ -330,12 +330,89 @@ export const getPath = async (
return isMatch ? [...parents, item] : parents; return isMatch ? [...parents, item] : parents;
}; };
// Polyfill for Promise.withResolvers
const Promise_withResolvers = <T>() => {
let resolve: (value: T | PromiseLike<T>) => void;
let reject: (reason?: unknown) => void;
const promise = new Promise<T>((res, rej) => {
resolve = res;
reject = rej;
});
return { promise, resolve: resolve!, reject: reject! };
};
// type ArchiverAppendArgs = Parameters<archiver.Archiver["append"]>;
type ArchiverAppendArgs = { readable: Readable; name: string; prefix: string };
/**
* The purpose of this class is to limit simultaneous open requests for readables.
* It supposes the list of file names and closures to get their readables fits in
* memory. Appending after being throttled is instantaneous.
*/
export class ThrottledArchive {
private readonly pendingAppends: (() => Promise<ArchiverAppendArgs>)[] = [];
private currentlyRunningAppends = 0;
private finalWaitingPromise?: PromiseWithResolvers<void>;
constructor(
public readonly actualArchive: archiver.Archiver,
private readonly appenderCb: (args: ArchiverAppendArgs) => Promise<void>,
private readonly maxConcurrentStorageReadsPerDownload: number = 4,
) {
if (maxConcurrentStorageReadsPerDownload < 1)
throw new Error(
`Invalid maxConcurrentStorageReadsPerDownload: ${JSON.stringify(
maxConcurrentStorageReadsPerDownload,
)}`,
);
}
private checkIsFinished() {
if (this.currentlyRunningAppends + this.pendingAppends.length === 0)
return this.finalWaitingPromise?.resolve();
}
private isAvailable(): boolean {
return this.currentlyRunningAppends < this.maxConcurrentStorageReadsPerDownload;
}
private async beginAnAppend(args: ArchiverAppendArgs) {
this.currentlyRunningAppends++;
args.readable.on("close", () => {
this.currentlyRunningAppends--;
this.maybeStartNextPending();
this.checkIsFinished();
});
await this.appenderCb(args);
}
private async maybeStartNextPending() {
if (!this.isAvailable() || this.pendingAppends.length < 1) return;
let args;
try {
args = await this.pendingAppends.pop()();
} catch (err) {
logger.error({ err }, "Error adding file to zip");
return;
}
return await this.beginAnAppend(args);
}
async append(getStreamCb: () => Promise<ArchiverAppendArgs>) {
this.pendingAppends.push(getStreamCb);
return this.maybeStartNextPending();
}
async waitUntilEmptied() {
if (this.finalWaitingPromise) {
return this.finalWaitingPromise.promise;
}
if (this.currentlyRunningAppends + this.pendingAppends.length === 0) {
return;
}
this.finalWaitingPromise = Promise_withResolvers();
return this.finalWaitingPromise.promise;
}
}
/** /**
* Adds drive items to an archive recursively * Adds drive items to an archive recursively
* *
* @param {string} id - the drive item id * @param {string} id - the drive item id
* @param {DriveFile | null } entity - the drive item entity * @param {DriveFile | null } entity - the drive item entity
* @param {archiver.Archiver} archive - the archive * @param {ThrottledArchive} throttledArchive - the archive to write to
* @param {Repository<DriveFile>} repository - the repository * @param {Repository<DriveFile>} repository - the repository
* @param {CompanyExecutionContext} context - the execution context * @param {CompanyExecutionContext} context - the execution context
* @param {string} prefix - folder prefix * @param {string} prefix - folder prefix
@@ -344,7 +421,7 @@ export const getPath = async (
export const addDriveItemToArchive = async ( export const addDriveItemToArchive = async (
id: string, id: string,
entity: DriveFile | null, entity: DriveFile | null,
archive: archiver.Archiver, throttledArchive: ThrottledArchive,
beginArchiveTransmit: beginArchiveTransmit:
| "ONLY_SET_THIS_VALUE_IF_ALREADY_CALLED" | "ONLY_SET_THIS_VALUE_IF_ALREADY_CALLED"
| ((archive: archiver.Archiver) => void), | ((archive: archiver.Archiver) => void),
@@ -362,6 +439,7 @@ export const addDriveItemToArchive = async (
if (item.is_in_trash) return; if (item.is_in_trash) return;
if (!item.is_directory) { if (!item.is_directory) {
await throttledArchive.append(async () => {
// random comment // random comment
const file_id = item.last_version_cache.file_metadata.external_id; const file_id = item.last_version_cache.file_metadata.external_id;
const file = await gr.services.files.download(file_id, context); const file = await gr.services.files.download(file_id, context);
@@ -369,13 +447,12 @@ export const addDriveItemToArchive = async (
if (!file) { if (!file) {
throw Error("file not found"); throw Error("file not found");
} }
return { readable: file.file, name: item.name, prefix: prefix ?? "" };
archive.append(file.file, { name: item.name, prefix: prefix ?? "" }); });
if (beginArchiveTransmit !== "ONLY_SET_THIS_VALUE_IF_ALREADY_CALLED") { if (beginArchiveTransmit !== "ONLY_SET_THIS_VALUE_IF_ALREADY_CALLED") {
beginArchiveTransmit(archive); beginArchiveTransmit(throttledArchive.actualArchive);
beginArchiveTransmit = "ONLY_SET_THIS_VALUE_IF_ALREADY_CALLED"; beginArchiveTransmit = "ONLY_SET_THIS_VALUE_IF_ALREADY_CALLED";
} }
return;
} else { } else {
let nextPage = ""; let nextPage = "";
do { do {
@@ -390,7 +467,7 @@ export const addDriveItemToArchive = async (
await addDriveItemToArchive( await addDriveItemToArchive(
child.id, child.id,
child, child,
archive, throttledArchive,
beginArchiveTransmit, beginArchiveTransmit,
repository, repository,
context, context,