mirror of
https://github.com/mollersuite/monofile.git
synced 2024-11-22 05:46:26 -08:00
readFileStream now respects backpressure more
This commit is contained in:
parent
7d18a60589
commit
fd8e143e3b
|
@ -106,9 +106,8 @@ async function startPushingWebStream(stream: Readable, webStream: ReadableStream
|
||||||
|
|
||||||
return reader.read().then(result => {
|
return reader.read().then(result => {
|
||||||
if (result.value)
|
if (result.value)
|
||||||
stream.push(result.value)
|
|
||||||
pushing = false
|
pushing = false
|
||||||
return result.done
|
return {readyForMore: result.value ? stream.push(result.value) : false, streamDone: result.done }
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
@ -544,7 +543,7 @@ export default class Files {
|
||||||
let d = await fetch(scanning_chunk.url, {headers})
|
let d = await fetch(scanning_chunk.url, {headers})
|
||||||
.catch((e: Error) => {
|
.catch((e: Error) => {
|
||||||
console.error(e)
|
console.error(e)
|
||||||
return {body: "__ERR"}
|
return {body: e}
|
||||||
})
|
})
|
||||||
|
|
||||||
position++
|
position++
|
||||||
|
@ -552,35 +551,50 @@ export default class Files {
|
||||||
return d.body
|
return d.body
|
||||||
}
|
}
|
||||||
|
|
||||||
let ord: number[] = []
|
let currentPusher : (() => Promise<{readyForMore: boolean, streamDone: boolean }> | undefined) | undefined
|
||||||
// hopefully this regulates it?
|
let busy = false
|
||||||
let lastChunkSent = true
|
|
||||||
|
|
||||||
let dataStream = new Readable({
|
let pushWS : (stream: Readable) => Promise<boolean | undefined> = async (stream: Readable) => {
|
||||||
read() {
|
|
||||||
if (!lastChunkSent) return
|
// uh oh, we don't have a currentPusher
|
||||||
lastChunkSent = false
|
// let's make one then
|
||||||
getNextChunk().then(async (nextChunk) => {
|
if (!currentPusher) {
|
||||||
if (typeof nextChunk == "string") {
|
let next = await getNextChunk()
|
||||||
this.destroy(new Error("file read error"))
|
if (next && !(next instanceof Error))
|
||||||
|
// okay, so we have a new chunk
|
||||||
|
// let's generate a new currentPusher
|
||||||
|
currentPusher = await startPushingWebStream(stream, next)
|
||||||
|
else {
|
||||||
|
// oops, look like there's an error
|
||||||
|
// or the stream has ended.
|
||||||
|
// let's destroy the stream
|
||||||
|
stream.destroy(next || undefined)
|
||||||
return
|
return
|
||||||
}
|
}
|
||||||
|
|
||||||
if (!nextChunk) return // EOF
|
|
||||||
|
|
||||||
let response = await pushWebStream(this, nextChunk)
|
|
||||||
|
|
||||||
|
|
||||||
while (response) {
|
|
||||||
let nextChunk = await getNextChunk()
|
|
||||||
// idk why this line was below but i moved it on top
|
|
||||||
// hopefully it wasn't for some other weird reason
|
|
||||||
if (!nextChunk || typeof nextChunk == "string") return
|
|
||||||
response = await pushWebStream(this, nextChunk)
|
|
||||||
}
|
}
|
||||||
lastChunkSent = true
|
|
||||||
})
|
let result = await currentPusher()
|
||||||
},
|
|
||||||
|
if (result?.streamDone) currentPusher = undefined;
|
||||||
|
return result?.readyForMore
|
||||||
|
|
||||||
|
}
|
||||||
|
|
||||||
|
let dataStream = new Readable({
|
||||||
|
async read() {
|
||||||
|
|
||||||
|
if (busy) return
|
||||||
|
busy = true
|
||||||
|
let readyForMore = true
|
||||||
|
|
||||||
|
while (readyForMore) {
|
||||||
|
let result = await pushWS(this)
|
||||||
|
if (result === undefined) return // stream has been destroyed. nothing left to do...
|
||||||
|
readyForMore = result
|
||||||
|
}
|
||||||
|
busy = false
|
||||||
|
|
||||||
|
}
|
||||||
})
|
})
|
||||||
|
|
||||||
return dataStream
|
return dataStream
|
||||||
|
|
Loading…
Reference in a new issue