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
70 changes: 32 additions & 38 deletions lib/multipart/parser.js
Original file line number Diff line number Diff line change
Expand Up @@ -345,53 +345,47 @@ const MultipartParser = class MultipartParser {
if (!this._sink) {
throw new Error("No sink configured");
}
const writer = await this._sink.write(assetPath, asset.mimeType, {
traceId,
});

const integrityStream = ssri.integrityStream({ single: true });
let hash = "";
integrityStream.once("integrity", (/** @type {any} */ integrity) => {
hash = integrity;
});

return new Promise((resolve, reject) => {
const onAbort = () => {
entry.destroy();
integrityStream.destroy();
writer.destroy();
};
if (signal.aborted) {
entry.resume();
throw new Error("Aborted");
}

// Buffer all chunks from the tar entry — no pipeline or stream Transform needed
const chunks = [];
for await (const chunk of entry) {
if (signal.aborted) {
onAbort();
reject(new Error("Aborted"));
return;
throw new Error("Aborted");
}
chunks.push(chunk);
}
const buffer = Buffer.concat(chunks);

signal.addEventListener("abort", onAbort, { once: true });

pipeline(entry, integrityStream, writer, (error) => {
signal.removeEventListener("abort", onAbort);

if (error) {
this._log.error(
`multipart - Failed writing asset to sink - Pathname: ${assetPath} - TraceId: ${traceId}`,
);
this._log.trace(error);
if (signal.aborted) {
throw new Error("Aborted");
}

reject(error);
return;
}
// Compute integrity hash directly from the buffer
const integrity = ssri.fromData(buffer, { single: true });
asset.integrity = integrity.toString();

asset.integrity = hash.toString();
try {
await this._sink.writeBuffer(assetPath, asset.mimeType, buffer, {
traceId,
});
} catch (error) {
this._log.error(
`multipart - Failed writing asset to sink - Pathname: ${assetPath} - TraceId: ${traceId}`,
);
this._log.trace(error);
throw error;
}

this._log.debug(
`multipart - Successfully wrote asset to sink - Pathname: ${assetPath} - TraceId: ${traceId}`,
);
this._log.debug(
`multipart - Successfully wrote asset to sink - Pathname: ${assetPath} - TraceId: ${traceId}`,
);

resolve(asset);
});
});
return asset;
}
};
export default MultipartParser;
57 changes: 57 additions & 0 deletions lib/sinks/test.js
Original file line number Diff line number Diff line change
Expand Up @@ -109,6 +109,63 @@ export default class SinkTest extends Sink {

// Common SINK API

/**
* @param {string} filePath
* @param {string} contentType
* @param {Buffer} buffer
* @param {{ traceId?: string }} [options]
* @returns {Promise<void>}
*/
// eslint-disable-next-line no-unused-vars
async writeBuffer(filePath, contentType, buffer, options = {}) {
const operation = "write";
try {
Sink.validateFilePath(filePath);
Sink.validateContentType(contentType);
} catch (error) {
this._counter.inc({ labels: { operation } });
throw error;
}
const pathname = toUrlPathname(path.join(this._rootPath, filePath));
if (pathname.indexOf(this._rootPath) !== 0) {
this._counter.inc({ labels: { operation } });
throw new Error(`Directory traversal - ${filePath}`);
}
const entry = new Entry({ mimeType: contentType, payload: [buffer] });
this._state.set(pathname, entry);
this._counter.inc({ labels: { success: true, access: true, operation } });
}

/**
* @param {string} filePath
* @returns {Promise<Buffer>}
*/
async readBuffer(filePath) {
const operation = "read";
try {
Sink.validateFilePath(filePath);
} catch (error) {
this._counter.inc({ labels: { operation } });
throw error;
}
const pathname = toUrlPathname(path.join(this._rootPath, filePath));
if (pathname.indexOf(this._rootPath) !== 0) {
this._counter.inc({ labels: { operation } });
throw new Error(`Directory traversal - ${filePath}`);
}
const entry = this._state.get(pathname);
if (!entry) {
this._counter.inc({ labels: { access: true, operation } });
throw new Error(`${filePath} does not exist`);
}
this._counter.inc({ labels: { success: true, access: true, operation } });
return Buffer.concat(
(entry.payload || []).map((/** @type {any} */ item) =>
Buffer.isBuffer(item) ? item : Buffer.from(item),
),
);
}

/**
* @param {string} filePath
* @param {string} contentType
Expand Down
63 changes: 10 additions & 53 deletions lib/utils/utils.js
Original file line number Diff line number Diff line change
@@ -1,39 +1,14 @@
import { Writable, Readable, pipeline } from "node:stream";
import { Writable, pipeline } from "node:stream";

/**
* @param {any} sink
* @param {string} path
*/
const readJSON = (sink, path) =>
// eslint-disable-next-line no-async-promise-executor
new Promise(async (resolve, reject) => {
try {
/** @type {any[]} */
const buffer = [];
const from = await sink.read(path);

const to = new Writable({
objectMode: false,
write(chunk, encoding, callback) {
buffer.push(chunk);
callback();
},
});
const readJSON = async (sink, path) => {
const buffer = await sink.readBuffer(path);
return JSON.parse(buffer.toString("utf8"));
};

pipeline(from.stream, to, (error) => {
if (error) return reject(error);
const str = Buffer.concat(buffer).toString("utf8");
try {
const obj = JSON.parse(str);
return resolve(obj);
} catch (err) {
return reject(err);
}
});
} catch (error) {
reject(error);
}
});
/**
* @param {any} sink
* @param {string} path
Expand All @@ -46,30 +21,12 @@ const readEikJson = (sink, path) => sink.exist(path);
* @param {any} obj
* @param {string} contentType
*/
const writeJSON = (sink, path, obj, contentType) =>
// eslint-disable-next-line no-async-promise-executor
new Promise(async (resolve, reject) => {
try {
const buffer = Buffer.from(JSON.stringify(obj));

const from = new Readable({
objectMode: false,
read() {
this.push(buffer);
this.push(null);
},
});

const to = await sink.write(path, contentType);
const writeJSON = async (sink, path, obj, contentType) => {
const buffer = Buffer.from(JSON.stringify(obj));
await sink.writeBuffer(path, contentType, buffer);
return buffer;
};

pipeline(from, to, (error) => {
if (error) return reject(error);
return resolve(buffer);
});
} catch (error) {
reject(error);
}
});
/**
* @param {any} from
*/
Expand Down
40 changes: 14 additions & 26 deletions package-lock.json

Some generated files are not rendered by default. Learn more about how customized files appear on GitHub.

6 changes: 3 additions & 3 deletions package.json
Original file line number Diff line number Diff line change
Expand Up @@ -29,9 +29,9 @@
"author": "",
"license": "MIT",
"dependencies": {
"@eik/sink": "1.3.0",
"@eik/sink-file-system": "2.3.0",
"@eik/sink-memory": "2.3.0",
"@eik/sink": "1.4.0",
"@eik/sink-file-system": "2.4.0",
"@eik/sink-memory": "2.4.0",
"@metrics/client": "2.5.5",
"abslog": "2.4.4",
"busboy": "1.6.0",
Expand Down
41 changes: 15 additions & 26 deletions test/multipart/parser.js
Original file line number Diff line number Diff line change
@@ -1,4 +1,4 @@
import { PassThrough, Writable } from "node:stream";
import { PassThrough } from "node:stream";
import HttpError from "http-errors";
import { URL } from "node:url";
import { test } from "node:test";
Expand Down Expand Up @@ -616,36 +616,25 @@ test("Parser() - _handleFile caps concurrent sink writes to avoid exhausting con
let activeWrites = 0;
let peakActiveWrites = 0;

// Simulates a real sink (e.g. GCS) where:
// - write() resolves immediately (just creates the upload stream object)
// - data drains through the stream without backpressure
// - but the stream stays "open" (final() is delayed) while data uploads
//
// This means multiple entries can be in-flight simultaneously, and the
// peak active count reflects real concurrent open upload streams.
// Simulates a real sink (e.g. GCS) where writeBuffer() is async and
// takes some time to complete, meaning multiple entries can be in-flight
// simultaneously. Peak active count reflects real concurrent uploads.
const trackingSink = {
write() {
writeBuffer() {
activeWrites++;
if (activeWrites > peakActiveWrites) {
peakActiveWrites = activeWrites;
}
return Promise.resolve(
new Writable({
write(chunk, _enc, cb) {
// Accept data instantly — no backpressure — so the tar parser
// can advance to the next entry before this one completes.
cb();
},
final(cb) {
// Delay completion to simulate an in-progress upload.
// During this window other entries can dispatch and open
// their own streams, making peak concurrency observable.
setTimeout(() => {
activeWrites--;
cb();
}, 20);
},
}),
// Delay completion to simulate an in-progress upload.
// During this window other entries can dispatch their own
// writeBuffer calls, making peak concurrency observable.
return /** @type {Promise<void>} */ (
new Promise((resolve) => {
setTimeout(() => {
activeWrites--;
resolve(undefined);
}, 20);
})
);
},
exist() {
Expand Down
Loading