Repository navigation
Summit Topic: Node.js Streams promise support #216
Description
Activity
- addedSession ProposalA session proposal for the Collaboration SummitA session proposal for the Collaboration Summit
on Dec 1, 2019 I’d wager this is probably “intermediate” and should list my session as a follow-up
@mcollina How is this different that the other session on Streams?
@WaleedAshraf - sounds good
@Fishrock123 - made note
@mcollina - How much time?From 30 minutes to 1 hour
Promisify
writable.write()Key Constraints:
highWaterMarksupport“Proper” Error handling
for await (let chunk of stream) { await ??? }Reacted by Trivikram Kamat and Matteo CollinaVery rough idea:
No promises writeable.write, instead we use duality and do the async iterator API.
Async iterators do:
{ next(value): Promise<nextValueResult> return(value): Promise<nextIfNoYieldOrFinallyOrLastValueResult> throw(err): Promise<nextValueResult> }
You can make a direct dual interface (moving parts around):
{ next(nextValueResult): Promise<value> return(nextIfNoYieldOrFinallyOrLastValueResult): Promise<value> throw(nextValueResult): MaybeRejectedPromsie<err> }
The first thing you will notice is that async iterators as an interface is already a dual interface. We haven't actually changed anything here in orxcer to get the API - since async iterators are bidiractional and kinda-dual already.
That is, if you had an API that did:
writeable.next(value)that returns a promise for the value and have that promise resolve fulfill when the value is ready that would be a nice API - you would need to support encoding and controlling whether or not it needs to wait for the highWatermark (since you have no drain). I would recommend always waiting for the highWatermark and if someone doesn't want that they can just set it to zero.Error handling would work like in async iterator.
Reacted by Robert NagyWhat do we put in HTTP.promises? nothing.
Another idea is service workers - but that is not really a server HTTP API.
I would recommend we just promisify the building blocks of HTTP (streams and event emitters).
For the client: probably
fetchonce we resolve eventemitter vs eventtarget etc.17 remaining items
@benjamingr to avoid unhandled promises when people do not actually need to wait for one or the other.
Other options:
- Regular promises on a "queued" object
- Events on a "queued" object
- People can join on the zoom link.On Fri, Dec 13, 2019 at 11:54 AM Ruben Bridgewater ***@***.***> wrote: Is it possible to stream this session by any chance? (As in open a zoom session) — You are receiving this because you were mentioned. Reply to this email directly, view it on GitHub <#216?email_source=notifications&email_token=AKYQ3QIVKXB736ZN2JLV3OTQYO45PA5CNFSM4JTN6RX2YY3PNVWWK3TUL52HS4DFVREXG43VMVBW63LNMVXHJKTDN5WW2ZLOORPWSZGOEG2QT3Y#issuecomment-565512687>, or unsubscribe <https://github.com/notifications/unsubscribe-auth/AKYQ3QKHQ6HNVOV46G7BYCTQYO45PANCNFSM4JTN6RXQ> .-- Eva Howe Operations Manager This Dot Labs eva@thisdot.co
@evahowe note that this session has ended a while ago and the next one is streaming :]
I feel like something like this works reasonably well, whether you care about errors or not:
Writable.prototype.write[promisify.custom] = () => function promisifiedWrite(chunk, encoding) { let syncErr; const needDrain = !this.write(chunk, encoding, (err) => { syncErr = err }); const stream = this; return { then() { if (syncErr) { throw syncErr; } if (needDrain) { return new Promise((resolve, reject) => { stream.once('error', reject); stream.once('drain', () => { stream.removeEventListener(reject); resolve(); }); }); } } }; };
Wouldn't leave it as a promisified function though, ideally it's first-class.
apapirovski did you mean to make that
thena getter? If so it needs to return thethenof tthe internal promise (and return a rejected promise and not throw an error). If not then the signature of then isthen(onFulfilled, onRejected)and both of them perform recursive assimilation. Any reason this can't be?Writable.prototype.write[promisify.custom] = () => function (chunk, encoding) { let syncErr; const needDrain = !this.write(chunk, encoding, (err) => { syncErr = err }); if (syncErr) return Promise.reject(err); // does this always happen synchronously? if (!needDrain) return Promise.resolve(); return EventEmitter.once(this, 'drain'); // once handles errors };
@benjamingr nope, meant that to be as is...
The implementation you've posted will reject a promise regardless. Avoiding that was an explicit design decision.
But also to be fair the throw is kinda sketchy. I don't think that might even work correctly lol. (Although seems to in my testing...) Either way, that line could be swapped for
return Promise.reject(syncErr)Anyway, sounds like you're suggesting something more like this...
Writable.prototype.write[promisify.custom] = () => function promisifiedWrite(chunk, encoding) { let syncErr; const needDrain = !this.write(chunk, encoding, (err) => { syncErr = err }); const stream = this; return { get then() { if (syncErr) { return Promise.reject(syncErr); } if (needDrain) { return new Promise((resolve, reject) => { stream.once('error', reject); stream.once('drain', () => { stream.removeEventListener(reject); resolve(); }); }); } } }; };
Similar in behavior but the downside is that inspecting the returned object in any way can cause the promises to be created.
Either way, my proposal technically can leak memory if someone awaits the promise too late and it required draining... Which I guess brings you to something more like this:
Writable.prototype.write[promisify.custom] = () => function promisifiedWrite(chunk, encoding) { let syncErr; const needDrain = !this.write(chunk, encoding, (err) => { syncErr = err }); let drainFn; if (needDrain) { stream.once('drain', () => { needDrain = false; if (drainFn) { drainFn(); } }); } const stream = this; return { then() { if (syncErr) { return Promise.reject(syncErr); } if (needDrain) { return new Promise((resolve, reject) => { stream.once('error', reject); drainFn = () => { stream.removeEventListener(reject); resolve(); }; }); } } }; };
The Promises/A+ spec (and thus the ECMAScript spec) explicitly requires calling
thenas a getter exactly once when chaining a promise though I see a few issues (like adding a listener on each call and also - probably return thethenof the returned promises bound to them rather than the promise itself (since then is a function and not a promise).I think lazy promises are mostly confusing (since promises are not actions) - and I agree that if we explore it we should not execute the action on the
thengetter but on invocation.
If what we want is to not create a rejection unless someone is listening and to not listen to certain events until someone is listening (which is arguably confusing) - maybe something like:
const lazy = fn => class extends Promise { // so we get catch, finally etc that all call `then` #promise = null; then(onFulfilled, onRejected) { this.#promise = this.#promise || fn(); return this.#promise.then(onFulfilled, onRejected); } }; Writable.prototype.write[promisify.custom] = () => function (chunk, encoding) { let syncErr; const needDrain = !this.write(chunk, encoding, (err) => { syncErr = err }); if (syncErr) { let writeError = Promise.reject(err); writeError.catch(() => {}); // don't care about this rejection no need to track it return writeError; } if (!needDrain) return Promise.resolve(); // no effects return lazy(() => EventEmitter.once(this, 'drain')); };
That's much neater haha
which is arguably confusing
Yeah, I was just taking off from where Jeremiah left it. I can see why some people would want to have the option to not get unhandled rejections in certain scenarios when using the promisified API.
(Also need to handle the
errorevent which I forgot to do, but yeah...)Tweet thread of the session https://twitter.com/trivikram/status/1205512468954521600
@mcollina are there any notes or other artifacts for this session? If so can you share them here? I am happy to make a PR and add them to the summit directory or feel free to raise one by yourself. Thanks!
Most of the discussion is captured in this issue. I plan to open a couple of issues on Node in the coming weeks to keep the discussion going.
Reacted by Christian BromannOk, gonna close this issue then.
This is my take:
Writable.prototype.write[promisify.custom] = () => function (chunk, encoding) { return new Promise((resolve, reject) => { let complete = false; const needDrain = !this.write(chunk, encoding, (err) => { complete = true; if (err) { reject(err) } else { resolve(); } this.off('error', reject); }); if (complete) { // Do nothing } else if (needDrain) { this.on('error', reject); } else { resolve(); } }); };
Of course all of these has the downside of:
let bytesSuccessfullyWritten = 0; await w.write(buf); bytesSuccessfullyWritten += buf.length; // BUG // Coming here does NOT mean that the above write was successful
I think an API like this would make more semantic sense:
let bytesSuccessfullyWritten = 0; let bytesBuffered = 0; for await (const chunks of source) { bufferedBytes += chunk.length; if (!w.write(chunk)) { await w.flush(); bytesSuccessfullyWritten += bytesBuffered; bytesBuffered = 0; } } await w.end();
Which would avoid confusion caused by the natural assumption that a successfully resolved promise from
writemeans that the write succeeded.I think the main problem with any variation of
Writable.prototype.write[promisify.custom]is that this means that each call towriteintroduces aPromiseeven when one isn't really necessary. Further, if this api is designed to return aPromise, it must systematically be awaited even when such waiting is unnecessary. I think the performance impact of such an API might make it a non-starter.I was thinking that some of the challenges around tracking which writes have been flushed and which haven't. Imagine scenarios in which there are multiple logical bits of logic 'concurrently' writing to the same
Writable. In that world, you might be more interested in whether 'your' chunks have been flushed or not and less so in the total, stateful 'flush state' of the `Writable.In that scenario, I imagine a
Promise-friendlyWriterobject whose role is to effect writes to aWritableand track their flush state. This means theWritercan be a fullyPromise-friendly api and no changes need to be made to theWritableinterface.
Hypothetical interface
interface Writer { /** * Create a Writer for a given Writable stream. * * _Note: No absolute need for this to be a class-style api. It could * easily be `require('stream').createWriter(writable)`._ */ new(writable: Stream.Writable): Writer; /** * Wait for the stream to be drained. * * Moral equivalent to `require('events').once(this.writable, 'drain')`. */ drained(): Promise<void>; /** * Wait for all chunks written via this `Writer` to be flushed. */ flushed(): Promise<void>; /** * Write a chunk of data to the stream and produce a `WriteResult`. * * @param chunk The chunk of data to be written * @param encoding Optional chunk encoding */ write(chunk: string | Buffer, encoding?: string): WriteResult; } interface WriteResult { /** * Return a Promise for when this chunk of data is flushed * * _Note: Could also be a getter. Important bit is that a promise is * only produced when this api is *used*._ */ flushed(): Promise<void>; /** * Indication of whether the consumer should block on the `Promise` returned by `Writer#drained` */ shouldWaitForDrain: boolean; }
Usage example:
const writer = new Writer(sink); for await (const chunk of source) { const result = writer.write(chunk); if (result.shouldWaitForDrain) { await writer.drained(); } // If I was interested in making sure each individual chunk was // flushed, I could uncomment the following: // // await result.flushed(); } // Now that all chunks from source have been written to sink, we want to // make sure that all those chunks have actually been flushed. // We're not interested in each specific chunk's flush state but instead // that _ALL_ chunks have been flushed here so we're not looking at // `WriteResult#flushed`. await writer.flushed();
Topic of the session
We have added async iterators to stream a quite successfully! It's time to think of what's missing.
Type of the session
Follow-up / Set-up sessions (if any)
Level
Pre-requisite knowledge
Describe the session
Session facilitator(s) and Github handle(s)
Additional context (optional)