Skip to content
File

Blob: src/workerd/api/streams/encoding.c++

8.6 KB
1// Copyright (c) 2017-2022 Cloudflare, Inc.
2// Licensed under the Apache 2.0 license found in the LICENSE file or at:
3// https://opensource.org/licenses/Apache-2.0
4 
5#include "encoding.h"
6 
7#include "simdutf.h"
8 
9#include <workerd/api/encoding.h>
10#include <workerd/api/streams/standard.h>
11#include <workerd/io/features.h>
12#include <workerd/jsg/jsg.h>
13 
14#include <v8.h>
15 
16#include <kj/common.h>
17#include <kj/refcount.h>
18 
19namespace workerd::api {
20 
21namespace {
22constexpr kj::byte REPLACEMENT_UTF8[] = {0xEF, 0xBF, 0xBD};
23 
24struct Holder: public kj::Refcounted {
25 kj::Maybe<char16_t> pending = kj::none;
26};
27} // namespace
28 
29// TextEncoderStream encodes a stream of JavaScript strings into UTF-8 bytes.
30//
31// WHATWG Encoding spec requirement (https://encoding.spec.whatwg.org/#interface-textencoderstream):
32// The encoder must encode unpaired UTF-16 surrogates as replacement characters.
33//
34// simdutf handles this for us, but we have to be careful of surrogate pairs
35// (high surrogate, followed by low surrogate) split across chunk boundaries.
36//
37// We do this with the pending field:
38// holder->pending = kj::none -> No pending high surrogate from previous chunk
39// holder->pending = char16_t -> High surrogate waiting for a matching low surrogate
40//
41// Ref: https://github.com/web-platform-tests/wpt/blob/master/encoding/streams/encode-utf8.any.js
42jsg::Ref<TextEncoderStream> TextEncoderStream::constructor(jsg::Lock& js) {
43 auto state = kj::rc<Holder>();
44 
45 auto transform = [holder = state.addRef()](jsg::Lock& js, v8::Local<v8::Value> chunk,
46 jsg::Ref<TransformStreamDefaultController> controller) mutable {
47 auto str = jsg::check(chunk->ToString(js.v8Context()));
48 size_t length = str->Length();
49 if (length == 0) return js.resolvedPromise();
50 
51 // Allocate buffer: reserve slot 0 for pending surrogate if we have one
52 size_t prefix = (holder->pending == kj::none) ? 0 : 1;
53 size_t end = prefix + length;
54 auto buf = kj::heapArray<char16_t>(end);
55 str->WriteV2(js.v8Isolate, 0, length, reinterpret_cast<uint16_t*>(buf.begin() + prefix));
56 
57 KJ_IF_SOME(lead, holder->pending) {
58 buf.begin()[0] = lead;
59 holder->pending = kj::none;
60 }
61 
62 // If chunk ends with high surrogate, save it for next chunk
63 if (end > 0 && U_IS_LEAD(buf[end - 1])) {
64 holder->pending = buf[--end];
65 }
66 if (end == 0) return js.resolvedPromise();
67 
68 auto slice = buf.first(end);
69 auto result = simdutf::utf8_length_from_utf16_with_replacement(slice.begin(), slice.size());
70 // Only sanitize if there are surrogates in the buffer - UTF-16 without
71 // surrogates is always well-formed.
72 if (result.error == simdutf::error_code::SURROGATE) {
73 simdutf::to_well_formed_utf16(slice.begin(), slice.size(), slice.begin());
74 }
75 auto utf8Length = result.count;
76 KJ_DASSERT(utf8Length > 0 && utf8Length >= end);
77 
78 auto backingStore = js.allocBackingStore(utf8Length, jsg::Lock::AllocOption::UNINITIALIZED);
79 auto dest = kj::ArrayPtr<char>(static_cast<char*>(backingStore->Data()), utf8Length);
80 [[maybe_unused]] auto written =
81 simdutf::convert_utf16_to_utf8(slice.begin(), slice.size(), dest.begin());
82 KJ_DASSERT(written == utf8Length, "simdutf should write exactly utf8Length bytes");
83 
84 auto array = v8::Uint8Array::New(
85 v8::ArrayBuffer::New(js.v8Isolate, kj::mv(backingStore)), 0, utf8Length);
86 controller->enqueue(js, jsg::JsUint8Array(array));
87 return js.resolvedPromise();
88 };
89 
90 auto flush = [holder = state.addRef()](
91 jsg::Lock& js, jsg::Ref<TransformStreamDefaultController> controller) mutable {
92 // If stream ends with orphaned high surrogate, emit replacement character
93 if (holder->pending != kj::none) {
94 auto backingStore = js.allocBackingStore(3, jsg::Lock::AllocOption::UNINITIALIZED);
95 memcpy(backingStore->Data(), REPLACEMENT_UTF8, 3);
96 controller->enqueue(js, jsg::JsUint8Array::create(js, kj::mv(backingStore), 0, 3));
97 }
98 return js.resolvedPromise();
99 };
100 
101 // Per the WHATWG Encoding spec, the readable side HWM should be 0, so writes
102 // block until a reader pulls. Previously StreamQueuingStrategy{} was passed,
103 // which bypassed the orDefault() in TransformStream::constructor and caused
104 // the readable HWM to default to 1, clearing backpressure at startup.
105 // Passing kj::none lets TransformStream apply the spec defaults (writable HWM=1,
106 // readable HWM=0).
107 kj::Maybe<StreamQueuingStrategy> readableStrategy;
108 if (!FeatureFlags::get(js).getEncoderStreamSpecCompliantBackpressure()) {
109 readableStrategy = StreamQueuingStrategy{};
110 }
111 auto transformer = TransformStream::constructor(js,
112 Transformer{.transform = jsg::Function<Transformer::TransformAlgorithm>(kj::mv(transform)),
113 .flush = jsg::Function<Transformer::FlushAlgorithm>(kj::mv(flush))},
114 StreamQueuingStrategy{}, kj::mv(readableStrategy));
115 
116 return js.alloc<TextEncoderStream>(transformer->getReadable(), transformer->getWritable());
117}
118 
119TextDecoderStream::TextDecoderStream(jsg::Ref<TextDecoder> decoder,
120 jsg::Ref<ReadableStream> readable,
121 jsg::Ref<WritableStream> writable)
122 : TransformStream(kj::mv(readable), kj::mv(writable)),
123 decoder(kj::mv(decoder)) {}
124 
125jsg::Ref<TextDecoderStream> TextDecoderStream::constructor(
126 jsg::Lock& js, jsg::Optional<kj::String> label, jsg::Optional<TextDecoderStreamInit> options) {
127 
128 auto decoder = TextDecoder::constructor(js, kj::mv(label), options.map([&js](auto& opts) {
129 return TextDecoder::ConstructorOptions{
130 // Previously this would default to true. The spec requires a default
131 // of false, however. When the pedanticWpt flag is not set, we continue
132 // to default as true.
133 .fatal = opts.fatal.orDefault(!FeatureFlags::get(js).getPedanticWpt()),
134 .ignoreBOM = opts.ignoreBOM.orDefault(false),
135 };
136 }));
137 
138 // The controller will store c++ references to both the readable and writable
139 // streams underlying controllers.
140 // See comment in TextEncoderStream::constructor for why we conditionally pass
141 // kj::none for the readable strategy.
142 kj::Maybe<StreamQueuingStrategy> readableStrategy;
143 if (!FeatureFlags::get(js).getEncoderStreamSpecCompliantBackpressure()) {
144 readableStrategy = StreamQueuingStrategy{};
145 }
146 auto transformer = TransformStream::constructor(js,
147 Transformer{.transform = jsg::Function<Transformer::TransformAlgorithm>( JSG_VISITABLE_LAMBDA(
148 (decoder = decoder.addRef()), (decoder),
149 (jsg::Lock& js, auto chunk, auto controller) {
150 JSG_REQUIRE(chunk->IsArrayBuffer() || chunk->IsArrayBufferView(), TypeError,
151 "This TransformStream is being used as a byte stream, "
152 "but received a value that is not a BufferSource.");
153 jsg::BufferSource source(js, chunk);
154 auto decoded =
155 JSG_REQUIRE_NONNULL(decoder->decodePtr(js, source.asArrayPtr(), false),
156 TypeError, "Failed to decode input.");
157 // Only enqueue if there's actual output - don't emit empty chunks
158 // for incomplete multi-byte sequences
159 if (decoded.length(js) > 0) {
160 controller->enqueue(js, decoded);
161 }
162 return js.resolvedPromise();
163 })),
164 .flush = jsg::Function<Transformer::FlushAlgorithm>(
165 JSG_VISITABLE_LAMBDA((decoder = decoder.addRef()), (decoder),
166 (jsg::Lock& js, auto controller) {
167 auto decoded =
168 JSG_REQUIRE_NONNULL(decoder->decodePtr(js, kj::ArrayPtr<kj::byte>(), true),
169 TypeError, "Failed to decode input.");
170 // Only enqueue if there's actual output
171 if (decoded.length(js) > 0) {
172 controller->enqueue(js, decoded);
173 }
174 return js.resolvedPromise();
175 }))},
176 StreamQueuingStrategy{}, kj::mv(readableStrategy));
177 
178 return js.alloc<TextDecoderStream>(
179 kj::mv(decoder), transformer->getReadable(), transformer->getWritable());
180}
181 
182kj::StringPtr TextDecoderStream::getEncoding() {
183 return decoder->getEncoding();
184}
185 
186bool TextDecoderStream::getFatal() {
187 return decoder->getFatal();
188}
189 
190bool TextDecoderStream::getIgnoreBOM() {
191 return decoder->getIgnoreBom();
192}
193 
194void TextDecoderStream::visitForMemoryInfo(jsg::MemoryTracker& tracker) const {
195 tracker.trackField("decoder", decoder);
196}
197 
198void TextDecoderStream::visitForGc(jsg::GcVisitor& visitor) {
199 visitor.visit(decoder);
200}
201 
202} // namespace workerd::api