Skip to main content

ferrijs_std/stream_web/
mod.rs

1use crate::utils::{
2    module::{export_default, ModuleInfo},
3    primordials::{BasePrimordials, Primordial},
4};
5use queuing_strategy::{ByteLengthQueuingStrategy, CountQueuingStrategy};
6use readable::{
7    ReadableByteStreamController, ReadableStreamBYOBReader, ReadableStreamBYOBRequest,
8    ReadableStreamDefaultController, ReadableStreamDefaultReader,
9};
10use rquickjs::{
11    module::{Declarations, Exports, ModuleDef},
12    Class, Ctx, Object, Result,
13};
14use writable::{WritableStream, WritableStreamDefaultController, WritableStreamDefaultWriter};
15
16use crate::stream_web::{
17    readable::{ArrayConstructorPrimordials, IteratorPrimordials},
18    transform::{TransformStream, TransformStreamDefaultController},
19    utils::promise::PromisePrimordials,
20    writable::WritableStreamDefaultControllerPrimordials,
21};
22
23mod queuing_strategy;
24pub mod readable;
25mod readable_writable_pair;
26mod transform;
27pub mod utils;
28mod writable;
29
30// Public API for creating streams from Rust
31pub use readable::stream::lock_readable_stream;
32pub use readable::stream::tee_readable_stream;
33pub use readable::stream::try_sync_drain_closed_stream;
34pub use readable::stream::ReadableStream;
35pub use readable::{
36    readable_byte_stream_controller_close_stream, readable_byte_stream_controller_enqueue_bytes,
37    readable_byte_stream_controller_enqueue_bytes_borrowed,
38    readable_stream_default_controller_close_stream,
39    readable_stream_default_controller_enqueue_value,
40    readable_stream_default_controller_error_stream, ReadableByteStreamControllerClass,
41    ReadableStreamDefaultControllerClass,
42};
43pub use readable::{CancelAlgorithm, PullAlgorithm, ReadableStreamControllerClass, StartAlgorithm};
44pub use readable::{NativePull, NativePullFn, NativePullResult};
45
46/// Creates a transform stream using LLRT's built-in Web Streams implementation.
47///
48/// This does not consult the global `TransformStream` binding.
49pub fn create_transform_stream<'js>(
50    ctx: &Ctx<'js>,
51    transformer: Object<'js>,
52) -> Result<Object<'js>> {
53    init_primordials(ctx)?;
54    Ok(TransformStream::from_transformer(ctx.clone(), transformer)?.into_inner())
55}
56
57fn init_primordials(ctx: &Ctx<'_>) -> Result<()> {
58    BasePrimordials::init(ctx)?;
59    PromisePrimordials::init(ctx)?;
60    ArrayConstructorPrimordials::init(ctx)?;
61    WritableStreamDefaultControllerPrimordials::init(ctx)?;
62    IteratorPrimordials::init(ctx)?;
63    Ok(())
64}
65
66/// Defines web streams, which are exposed through the "stream/web" Node import, but also at the global scope
67/// Web streams consist of Readable, Writable, and Transform streams. Transform is currently unimplemented.
68///
69/// https://developer.mozilla.org/en-US/docs/Web/API/Streams_API
70///
71/// # ReadableStream
72/// ReadableStream knows how to 'pull' objects or bytes from an underlying source, generally a user-defined function or an [async] iterator.
73/// A source enqueues data to the stream via a controller, either ReadableStreamDefaultController or a ReadableByteStreamController optionally for byte data.
74/// The controller is created at stream initialisation and cannot change.
75///
76/// Data is read from the stream using a reader, which is obtained using stream.getReader(). A reader 'locks' the stream for reading, preventing
77/// other readers from being created. When a reader is released with `reader.releaseLock()`, the stream goes back to having no reader and a new one can be created.
78/// In the case of ReadableByteStreamController, a special reader ReadableStreamBYOBReader may be used, which allows users to provide their own
79/// buffer to fill bytes into when reading. Otherwise, ReadableStreamDefaultReader is used by default, and this may also be used with byte streams.
80///
81/// A ReadableStream can be 'tee'd', which splits it into two readable streams which both read the same underlying data, potentially at different
82/// paces. This is an area of substantial complexity for the implementation, particularly in the case of byte streams as the alternative reader types
83/// must be handled correctly.
84///
85/// # WritableStream
86/// WritableStream knows how to 'push' objects into an underlying sink, generally a user-defined function. It has no special casing for bytes, and so
87/// only has one type of controller, WritableStreamDefaultController, and only one type of writer WritableStreamDefaultWriter. The controller is only needed for
88/// error handling because writes are signalled via a function call to a user-defined 'write' method which receives the chunk directly.
89///
90/// Data is written to the stream using a WritableStreamDefaultWriter, which is obtained using stream.getWriter(). A writer 'locks' the stream for writing,
91/// preventing other writers from being created. When a writer is released with `writer.releaseLock()`, the stream goes back to having no writer and a new one can be created.
92pub struct StreamWebModule;
93
94// https://nodejs.org/api/webstreams.html
95impl ModuleDef for StreamWebModule {
96    fn declare(declare: &Declarations) -> Result<()> {
97        declare.declare(stringify!(ReadableStream))?;
98        declare.declare(stringify!(ReadableStreamDefaultReader))?;
99        declare.declare(stringify!(ReadableStreamBYOBReader))?;
100        declare.declare(stringify!(ReadableStreamDefaultController))?;
101        declare.declare(stringify!(ReadableByteStreamController))?;
102        declare.declare(stringify!(ReadableStreamBYOBRequest))?;
103
104        declare.declare(stringify!(WritableStream))?;
105        declare.declare(stringify!(WritableStreamDefaultWriter))?;
106        declare.declare(stringify!(WritableStreamDefaultController))?;
107
108        declare.declare(stringify!(TransformStream))?;
109        declare.declare(stringify!(TransformStreamDefaultController))?;
110
111        declare.declare(stringify!(ByteLengthQueuingStrategy))?;
112        declare.declare(stringify!(CountQueuingStrategy))?;
113
114        declare.declare("default")?;
115        Ok(())
116    }
117
118    #[inline]
119    fn evaluate<'js>(ctx: &Ctx<'js>, exports: &Exports<'js>) -> Result<()> {
120        export_default(ctx, exports, |default| {
121            Class::<ReadableStream>::define(default)?;
122            Class::<ReadableStreamDefaultReader>::define(default)?;
123            Class::<ReadableStreamBYOBReader>::define(default)?;
124            Class::<ReadableStreamDefaultController>::define(default)?;
125            Class::<ReadableByteStreamController>::define(default)?;
126            Class::<ReadableStreamBYOBRequest>::define(default)?;
127
128            Class::<WritableStream>::define(default)?;
129            Class::<WritableStreamDefaultWriter>::define(default)?;
130            Class::<WritableStreamDefaultController>::define(default)?;
131
132            Class::<ByteLengthQueuingStrategy>::define(default)?;
133            Class::<CountQueuingStrategy>::define(default)?;
134
135            Class::<TransformStream>::define(default)?;
136            Class::<TransformStreamDefaultController>::define(default)?;
137
138            Ok(())
139        })?;
140
141        Ok(())
142    }
143}
144
145impl From<StreamWebModule> for ModuleInfo<StreamWebModule> {
146    fn from(val: StreamWebModule) -> Self {
147        ModuleInfo {
148            name: "stream/web",
149            module: val,
150        }
151    }
152}
153
154pub fn init(ctx: &Ctx) -> Result<()> {
155    let globals = &ctx.globals();
156
157    init_primordials(ctx)?;
158
159    // https://min-common-api.proposal.wintertc.org/#api-index
160    Class::<ByteLengthQueuingStrategy>::define(globals)?;
161    Class::<CountQueuingStrategy>::define(globals)?;
162
163    Class::<ReadableByteStreamController>::define(globals)?;
164    Class::<ReadableStream>::define(globals)?;
165    Class::<ReadableStreamBYOBReader>::define(globals)?;
166    Class::<ReadableStreamBYOBRequest>::define(globals)?;
167    Class::<ReadableStreamDefaultController>::define(globals)?;
168    Class::<ReadableStreamDefaultReader>::define(globals)?;
169
170    Class::<WritableStream>::define(globals)?;
171    Class::<WritableStreamDefaultController>::define(globals)?;
172
173    // This is exposed globally by Node even though its not in the min-common-api
174    Class::<WritableStreamDefaultWriter>::define(globals)?;
175
176    Class::<TransformStream>::define(globals)?;
177    Class::<TransformStreamDefaultController>::define(globals)?;
178
179    Ok(())
180}