|
|
@@ -23,6 +23,7 @@ class ObserverStream extends Stream.Transform {
|
|
|
if (options.hash) {
|
|
|
this.hash = crypto.createHash(options.hash)
|
|
|
}
|
|
|
+
|
|
|
if (options.metric) {
|
|
|
const onEnd = () => {
|
|
|
metrics.count(options.metric, this.bytes)
|
|
|
@@ -98,35 +99,61 @@ async function verifyMd5(persistor, bucket, key, sourceMd5, destMd5 = null) {
|
|
|
function getReadyPipeline(...streams) {
|
|
|
return new Promise((resolve, reject) => {
|
|
|
const lastStream = streams.slice(-1)[0]
|
|
|
- let resolvedOrErrored = false
|
|
|
|
|
|
+ // in case of error or stream close, we must ensure that we drain the
|
|
|
+ // previous stream so that it can clean up its socket (if it has one)
|
|
|
+ const drainPreviousStream = function(previousStream) {
|
|
|
+ // this stream is no longer reliable, so don't pipe anything more into it
|
|
|
+ previousStream.unpipe(this)
|
|
|
+ previousStream.resume()
|
|
|
+ }
|
|
|
+
|
|
|
+ // handler to resolve when either:
|
|
|
+ // - an error happens, or
|
|
|
+ // - the last stream in the chain is readable
|
|
|
+ // for example, in the case of a 4xx error an error will occur and the
|
|
|
+ // streams will not become readable
|
|
|
const handler = function(err) {
|
|
|
- if (!resolvedOrErrored) {
|
|
|
- resolvedOrErrored = true
|
|
|
-
|
|
|
- lastStream.removeListener('readable', handler)
|
|
|
- if (err) {
|
|
|
- reject(
|
|
|
- wrapError(err, 'error before stream became ready', {}, ReadError)
|
|
|
- )
|
|
|
- } else {
|
|
|
- resolve(lastStream)
|
|
|
- }
|
|
|
+ // remove handler from all streams because we don't want to do this on
|
|
|
+ // later errors
|
|
|
+ lastStream.removeListener('readable', handler)
|
|
|
+ for (const stream of streams) {
|
|
|
+ stream.removeListener('error', handler)
|
|
|
}
|
|
|
+
|
|
|
+ // return control to the caller
|
|
|
if (err) {
|
|
|
- for (const stream of streams) {
|
|
|
- if (!stream.destroyed) {
|
|
|
- stream.destroy()
|
|
|
- }
|
|
|
- }
|
|
|
+ reject(
|
|
|
+ wrapError(err, 'error before stream became ready', {}, ReadError)
|
|
|
+ )
|
|
|
+ } else {
|
|
|
+ resolve(lastStream)
|
|
|
}
|
|
|
}
|
|
|
|
|
|
+ // ensure the handler fires when the last strem becomes readable
|
|
|
+ lastStream.on('readable', handler)
|
|
|
+
|
|
|
+ for (const stream of streams) {
|
|
|
+ // when a stream receives a pipe, set up the drain handler to drain the
|
|
|
+ // connection if an error occurs or the stream is closed
|
|
|
+ stream.on('pipe', previousStream => {
|
|
|
+ stream.on('error', x => {
|
|
|
+ drainPreviousStream(previousStream)
|
|
|
+ })
|
|
|
+ stream.on('close', () => {
|
|
|
+ drainPreviousStream(previousStream)
|
|
|
+ })
|
|
|
+ })
|
|
|
+ // add the handler function to resolve this method on error if we can't
|
|
|
+ // set up the pipeline
|
|
|
+ stream.on('error', handler)
|
|
|
+ }
|
|
|
+
|
|
|
+ // begin the pipeline
|
|
|
for (let index = 0; index < streams.length - 1; index++) {
|
|
|
- streams[index + 1].on('close', () => streams[index].destroy())
|
|
|
+ streams[index].pipe(streams[index + 1])
|
|
|
}
|
|
|
- pipeline(...streams).catch(handler)
|
|
|
- lastStream.on('readable', handler)
|
|
|
})
|
|
|
}
|
|
|
|