diff --git a/experiments/WebScene.NativeEngine.Probe/native/webscene_v8_runtime.cpp b/experiments/WebScene.NativeEngine.Probe/native/webscene_v8_runtime.cpp index d6c112192..5c8929b2e 100644 --- a/experiments/WebScene.NativeEngine.Probe/native/webscene_v8_runtime.cpp +++ b/experiments/WebScene.NativeEngine.Probe/native/webscene_v8_runtime.cpp @@ -4431,6 +4431,176 @@ struct v8_dom_runtime::implementation final { } } + const writableStreamState = new WeakMap(); + const writableStreamWriterState = new WeakMap(); + const requireWritableStream = value => { + const state = writableStreamState.get(value); + if (!state) throw new TypeError('Illegal invocation'); + return state; + }; + const requireWritableStreamWriter = value => { + const state = writableStreamWriterState.get(value); + if (!state) throw new TypeError('Illegal invocation'); + return state; + }; + class WritableStreamDefaultWriter { + constructor(stream) { + const state = requireWritableStream(stream); + if (state.locked) throw new TypeError('WritableStream is locked'); + state.locked = true; + writableStreamWriterState.set(this, {stream, state}); + } + get closed() { return requireWritableStreamWriter(this).state.closed; } + get ready() { return requireWritableStreamWriter(this).state.ready; } + get desiredSize() { + const state = requireWritableStreamWriter(this).state; + return state.status === 'errored' ? null + : state.status === 'closed' ? 0 : 1; + } + write(chunk) { + const {stream, state} = requireWritableStreamWriter(this); + if (!stream) return Promise.reject(new TypeError('Writer released')); + if (state.status !== 'writable') { + return Promise.reject(state.error ?? new TypeError('WritableStream is closed')); + } + const operation = state.chain.then(() => state.write?.(chunk)); + state.chain = operation.catch(error => { + state.status = 'errored'; state.error = error; + state.closedReject(error); throw error; + }); + state.chain.catch(() => {}); + return operation; + } + close() { + const {stream, state} = requireWritableStreamWriter(this); + if (!stream) return Promise.reject(new TypeError('Writer released')); + if (state.status !== 'writable') { + return Promise.reject(state.error ?? new TypeError('WritableStream is closed')); + } + state.status = 'closing'; + const operation = state.chain.then(() => state.close?.()); + state.chain = operation.then(() => { + state.status = 'closed'; state.closedResolve(); + }, error => { + state.status = 'errored'; state.error = error; + state.closedReject(error); throw error; + }); + state.chain.catch(() => {}); + return operation; + } + abort(reason = undefined) { + const {stream, state} = requireWritableStreamWriter(this); + if (!stream) return Promise.reject(new TypeError('Writer released')); + if (state.status === 'closed') return Promise.resolve(); + state.status = 'errored'; state.error = reason; + state.closedReject(reason); + try { return Promise.resolve(state.abort?.(reason)); } + catch (error) { return Promise.reject(error); } + } + releaseLock() { + const value = requireWritableStreamWriter(this); + if (!value.stream) return; + value.state.locked = false; + value.stream = undefined; + } + } + class WritableStream { + constructor(underlyingSink = {}) { + if (underlyingSink === null || typeof underlyingSink !== 'object') { + throw new TypeError('underlyingSink must be an object'); + } + let closedResolve; + let closedReject; + const closed = new Promise((resolve, reject) => { + closedResolve = resolve; closedReject = reject; + }); + closed.catch(() => {}); + const state = { + locked:false, status:'writable', error:undefined, + write:typeof underlyingSink.write === 'function' + ? underlyingSink.write.bind(underlyingSink) : undefined, + close:typeof underlyingSink.close === 'function' + ? underlyingSink.close.bind(underlyingSink) : undefined, + abort:typeof underlyingSink.abort === 'function' + ? underlyingSink.abort.bind(underlyingSink) : undefined, + chain:Promise.resolve(), ready:Promise.resolve(), +)JS", + R"JS( closed, closedResolve, closedReject + }; + writableStreamState.set(this, state); + if (typeof underlyingSink.start === 'function') { + try { state.chain = Promise.resolve(underlyingSink.start()); } + catch (error) { state.chain = Promise.reject(error); } + } + state.chain.catch(error => { + state.status = 'errored'; state.error = error; state.closedReject(error); + }); + } + get locked() { return requireWritableStream(this).locked; } + getWriter() { return new WritableStreamDefaultWriter(this); } + abort(reason = undefined) { + const state = requireWritableStream(this); + if (state.locked) return Promise.reject(new TypeError('WritableStream is locked')); + const writer = new WritableStreamDefaultWriter(this); + const result = writer.abort(reason); + writer.releaseLock(); + return result; + } + close() { + const state = requireWritableStream(this); + if (state.locked) return Promise.reject(new TypeError('WritableStream is locked')); + const writer = new WritableStreamDefaultWriter(this); + const result = writer.close(); + writer.releaseLock(); + return result; + } + } + const transformStreamState = new WeakMap(); + class TransformStream { + constructor(transformer = {}) { + let controller; + const readable = new ReadableStream({start(value) { controller = value; }}); + const transform = typeof transformer.transform === 'function' + ? transformer.transform.bind(transformer) : undefined; + const flush = typeof transformer.flush === 'function' + ? transformer.flush.bind(transformer) : undefined; + const writable = new WritableStream({ + async write(chunk) { + if (transform) await transform(chunk, controller); + else controller.enqueue(chunk); + }, + async close() { if (flush) await flush(controller); controller.close(); }, + abort(reason) { controller.error(reason); } + }); + transformStreamState.set(this, {readable, writable}); + } + get readable() { + const state = transformStreamState.get(this); + if (!state) throw new TypeError('Illegal invocation'); + return state.readable; + } + get writable() { + const state = transformStreamState.get(this); + if (!state) throw new TypeError('Illegal invocation'); + return state.writable; + } + } + for (const [prototype, tag, names] of [ + [WritableStream.prototype, 'WritableStream', + ['locked', 'getWriter', 'abort', 'close']], + [WritableStreamDefaultWriter.prototype, 'WritableStreamDefaultWriter', + ['closed', 'ready', 'desiredSize', 'write', 'close', 'abort', 'releaseLock']], + [TransformStream.prototype, 'TransformStream', ['readable', 'writable']] + ]) { + Object.defineProperty(prototype, Symbol.toStringTag, + {value:tag, configurable:true}); + for (const name of names) { + const descriptor = Object.getOwnPropertyDescriptor(prototype, name); + if (descriptor) Object.defineProperty( + prototype, name, {...descriptor, enumerable:true}); + } + } + const bodyBytes = body => body instanceof ArrayBuffer ? new Uint8Array(body.slice(0)) : ArrayBuffer.isView(body) @@ -4708,6 +4878,15 @@ struct v8_dom_runtime::implementation final { ReadableStreamDefaultController: { value: ReadableStreamDefaultController, writable: true, configurable: true }, + WritableStream: { + value: WritableStream, writable: true, configurable: true + }, + WritableStreamDefaultWriter: { + value: WritableStreamDefaultWriter, writable: true, configurable: true + }, + TransformStream: { + value: TransformStream, writable: true, configurable: true + }, Request: { value: WebSceneRequest, writable: true, configurable: true }, @@ -4724,8 +4903,10 @@ struct v8_dom_runtime::implementation final { })(); )JS", }; + size_t source_size = 0U; + for (const auto part : source_parts) source_size += part.size(); std::string source; - source.reserve(27687); + source.reserve(source_size); for (const auto part : source_parts) source.append(part); auto script = v8::Script::Compile( local_context, diff --git a/experiments/WebScene.NativeEngine.Probe/tests/native_v8_runtime_stream_fetch_tests.inc b/experiments/WebScene.NativeEngine.Probe/tests/native_v8_runtime_stream_fetch_tests.inc index 1c685c299..2f4dd315a 100644 --- a/experiments/WebScene.NativeEngine.Probe/tests/native_v8_runtime_stream_fetch_tests.inc +++ b/experiments/WebScene.NativeEngine.Probe/tests/native_v8_runtime_stream_fetch_tests.inc @@ -250,11 +250,17 @@ void test_cache_storage_and_controlled_fetch_broker() const known = await cache.match(event.request); if (known) return known; ++misses; - const response = new Response('broker-ok', { + const transform = new TransformStream(); + const writer = transform.writable.getWriter(); + const response = new Response(transform.readable, { status:200, headers:{'content-type':'text/plain', 'x-resource-plane':'native'} }); const stored = cache.put(event.request, response.clone()); + await writer.write(new Uint8Array([98, 114, 111, 107, 101, 114, 45])); + await writer.write(new Uint8Array([111, 107])); + await writer.close(); + await writer.closed; await stored; return response; })()); @@ -266,7 +272,13 @@ void test_cache_storage_and_controlled_fetch_broker() cachesTag:Object.prototype.toString.call(caches), cacheStorageBrand:caches instanceof CacheStorage, cacheMethodEnumerable:Object.getOwnPropertyDescriptor( - CacheStorage.prototype, 'open').enumerable + CacheStorage.prototype, 'open').enumerable, + transformTag:Object.prototype.toString.call(new TransformStream()), + transformReadable:Object.getOwnPropertyDescriptor( + TransformStream.prototype, 'readable').enumerable, + writableTag:Object.prototype.toString.call(new WritableStream()), + writerMethodEnumerable:Object.getOwnPropertyDescriptor( + WritableStreamDefaultWriter.prototype, 'write').enumerable }); }); )JS"} @@ -297,6 +309,10 @@ void test_cache_storage_and_controlled_fetch_broker() }; (async () => { try { + if (Object.prototype.toString.call(new TransformStream()) !== + '[object TransformStream]') { + throw new Error('top-level TransformStream brand'); + } const changed = new Promise(resolve => navigator.serviceWorker.addEventListener( 'controllerchange', resolve, {once:true})); const registration = await navigator.serviceWorker.register( @@ -353,9 +369,13 @@ void test_cache_storage_and_controlled_fetch_broker() const value = __resourcePlane.report; return value.fetches === 100 && value.misses === 1 && value.emptyMatch && value.cachesTag === '[object CacheStorage]' && value.cacheStorageBrand - && value.cacheMethodEnumerable; + && value.cacheMethodEnumerable + && value.transformTag === '[object TransformStream]' + && value.transformReadable + && value.writableTag === '[object WritableStream]' + && value.writerMethodEnumerable; })())JS", "resource-plane-report.js") == "true", - "CacheStorage/controlled broker worker contract failed: " + "CacheStorage/TransformStream worker contract failed: " + evaluate(engine, "JSON.stringify(__resourcePlane.report)", "resource-plane-report-detail.js")); const auto p95 = std::stod(evaluate(engine, R"JS((() => { diff --git a/tests/WebPlatformSubset/contracts/cache-storage-transform-stream-worker.js b/tests/WebPlatformSubset/contracts/cache-storage-transform-stream-worker.js new file mode 100644 index 000000000..b570e3a82 --- /dev/null +++ b/tests/WebPlatformSubset/contracts/cache-storage-transform-stream-worker.js @@ -0,0 +1,55 @@ +let fetches = 0; +let misses = 0; +let emptyMatch = false; + +self.addEventListener('install', event => event.waitUntil(self.skipWaiting())); +self.addEventListener('activate', event => event.waitUntil((async () => { + const cache = await caches.open('cache-transform-v1'); + emptyMatch = await cache.match( + new Request(`${self.location.origin}/never`)) === undefined; + await self.clients.claim(); +})())); +self.addEventListener('fetch', event => { + if (!event.request.url.endsWith('/cache-transform-oracle')) return; + event.respondWith((async () => { + ++fetches; + const cache = await caches.open('cache-transform-v1'); + const known = await cache.match(event.request); + if (known) return known; + ++misses; + const transform = new TransformStream(); + const writer = transform.writable.getWriter(); + const response = new Response(transform.readable, { + status:200, headers:{'content-type':'text/plain'} + }); + const stored = cache.put(event.request, response.clone()); + await writer.write(new Uint8Array([98, 114, 111, 107, 101, 114, 45])); + await writer.write(new Uint8Array([111, 107])); + await writer.close(); + await writer.closed; + await stored; + return response; + })()); +}); +self.addEventListener('message', async event => { + if (event.data !== 'probe') return; + const cache = await caches.open('cache-transform-v1'); + const transform = new TransformStream(); + const writer = transform.writable.getWriter(); + event.ports[0].postMessage({ + fetches, misses, emptyMatch, + cachesTag:Object.prototype.toString.call(caches), + cacheStorageBrand:caches instanceof CacheStorage, + cacheTag:Object.prototype.toString.call(cache), + cacheMethodEnumerable:Object.getOwnPropertyDescriptor( + CacheStorage.prototype, 'open').enumerable, + transformTag:Object.prototype.toString.call(transform), + writableTag:Object.prototype.toString.call(transform.writable), + writerTag:Object.prototype.toString.call(writer), + transformReadable:Object.getOwnPropertyDescriptor( + TransformStream.prototype, 'readable').enumerable, + writerMethodEnumerable:Object.getOwnPropertyDescriptor( + WritableStreamDefaultWriter.prototype, 'write').enumerable + }); + writer.releaseLock(); +}); diff --git a/tests/WebPlatformSubset/contracts/cache-storage-transform-stream.html b/tests/WebPlatformSubset/contracts/cache-storage-transform-stream.html new file mode 100644 index 000000000..5754ec663 --- /dev/null +++ b/tests/WebPlatformSubset/contracts/cache-storage-transform-stream.html @@ -0,0 +1,53 @@ + + +