From c6c266498a0b068392324bb3b240f90dc763e054 Mon Sep 17 00:00:00 2001 From: =?UTF-8?q?Wies=C5=82aw=20=C5=A0olt=C3=A9s?= Date: Thu, 17 Sep 2026 19:55:43 +0200 Subject: [PATCH] Add TransformStream resource fallback --- .../native/webscene_v8_runtime.cpp | 183 +++++++++++++++++- .../native_v8_runtime_stream_fetch_tests.inc | 28 ++- .../cache-storage-transform-stream-worker.js | 55 ++++++ .../cache-storage-transform-stream.html | 53 +++++ .../webscene-component-profile.json | 8 + 5 files changed, 322 insertions(+), 5 deletions(-) create mode 100644 tests/WebPlatformSubset/contracts/cache-storage-transform-stream-worker.js create mode 100644 tests/WebPlatformSubset/contracts/cache-storage-transform-stream.html 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 @@ + + +CacheStorage, TransformStream, and controlled fetch broker + diff --git a/tests/WebPlatformSubset/webscene-component-profile.json b/tests/WebPlatformSubset/webscene-component-profile.json index 1a5a064e6..ef09a2817 100644 --- a/tests/WebPlatformSubset/webscene-component-profile.json +++ b/tests/WebPlatformSubset/webscene-component-profile.json @@ -2404,6 +2404,14 @@ "evidence": ["service-workers-cache-derived", "chromium-153-direct-neutral-contract", "vscode-oss-resource-broker-shape"], "reason": "Candidate second resource-plane stack for WebScene #266. It provides bounded worker-local CacheStorage with empty-match and cloned-response semantics plus asynchronous controlled Window.fetch dispatch with unhandled host fallback and rejected-response propagation. The native companion gate covers 100 controlled requests, one cache miss followed by stable clones, WebIDL brands/descriptors, unregister teardown, p95 <= 250 ms, <= 8 MiB retained V8 heap, and <= 64 MiB peak-RSS growth. WritableStream/TransformStream fallback, interception of parser/image/style loads, and unchanged Markdown product acceptance remain later #266 slices.", "nativeNavigation": true + }, + { + "path": "contracts/cache-storage-transform-stream.html", + "type": "contract", + "capabilities": ["transform-stream-fallback", "writable-stream-default-writer", "cache-transform-response-clone"], + "evidence": ["streams-transform-stream-derived", "chromium-153-direct-neutral-contract", "vscode-oss-safari-resource-stream-shape"], + "reason": "Candidate third resource-plane stack for WebScene #266. It adds WritableStream/default-writer and TransformStream byte delivery to the lower CacheStorage/request-broker stack and covers the unchanged Code OSS chunk fallback sequence: create a transform, retain its writer, resolve the response with its readable side, write ordered byte chunks, close or abort, and observe writer.closed teardown. The cumulative native gate covers 100 controlled cached requests, exact byte order, browser-shaped brands/descriptors, p95 <= 250 ms, <= 8 MiB retained V8 heap, and <= 64 MiB peak-RSS growth. Interception of parser/image/style loads and unchanged Markdown product acceptance remain later #266 slices.", + "nativeNavigation": true } ], "harnessBlocked": [],