simdjson 5.0.2
Ridiculously Fast JSON
Loading...
Searching...
No Matches
document_stream-inl.h
1#ifndef SIMDJSON_DOCUMENT_STREAM_INL_H
2#define SIMDJSON_DOCUMENT_STREAM_INL_H
3
4#include "simdjson/dom/base.h"
5#include "simdjson/dom/document_stream.h"
6#include "simdjson/dom/element-inl.h"
7#include "simdjson/dom/parser-inl.h"
8#include "simdjson/error-inl.h"
9#include "simdjson/internal/dom_parser_implementation.h"
10
11namespace simdjson {
12namespace dom {
13
14#ifdef SIMDJSON_THREADS_ENABLED
15
16inline void stage1_worker::finish() {
17 // After calling "run" someone would call finish() to wait
18 // for the end of the processing.
19 // This function will wait until either the thread has done
20 // the processing or, else, the destructor has been called.
21 std::unique_lock<std::mutex> lock(locking_mutex);
22 cond_var.wait(lock, [this]{return has_work == false;});
23}
24
25inline stage1_worker::~stage1_worker() {
26 // The thread may never outlive the stage1_worker instance
27 // and will always be stopped/joined before the stage1_worker
28 // instance is gone.
29 stop_thread();
30}
31
32inline void stage1_worker::start_thread() {
33 std::unique_lock<std::mutex> lock(locking_mutex);
34 if(thread.joinable()) {
35 return; // This should never happen but we never want to create more than one thread.
36 }
37 thread = std::thread([this]{
38 while(true) {
39 std::unique_lock<std::mutex> thread_lock(locking_mutex);
40 // We wait for either "run" or "stop_thread" to be called.
41 cond_var.wait(thread_lock, [this]{return has_work || !can_work;});
42 // If, for some reason, the stop_thread() method was called (i.e., the
43 // destructor of stage1_worker is called, then we want to immediately destroy
44 // the thread (and not do any more processing).
45 if(!can_work) {
46 break;
47 }
48 this->owner->stage1_thread_error = this->owner->run_stage1(*this->stage1_thread_parser,
49 this->_next_batch_start);
50 this->has_work = false;
51 // The condition variable call should be moved after thread_lock.unlock() for performance
52 // reasons but thread sanitizers may report it as a data race if we do.
53 // See https://stackoverflow.com/questions/35775501/c-should-condition-variable-be-notified-under-lock
54 cond_var.notify_one(); // will notify "finish"
55 thread_lock.unlock();
56 }
57 }
58 );
59}
60
61
62inline void stage1_worker::stop_thread() {
63 std::unique_lock<std::mutex> lock(locking_mutex);
64 // We have to make sure that all locks can be released.
65 can_work = false;
66 has_work = false;
67 cond_var.notify_all();
68 lock.unlock();
69 if(thread.joinable()) {
70 thread.join();
71 }
72}
73
74inline void stage1_worker::run(document_stream * ds, dom::parser * stage1, size_t next_batch_start) {
75 std::unique_lock<std::mutex> lock(locking_mutex);
76 owner = ds;
77 _next_batch_start = next_batch_start;
78 stage1_thread_parser = stage1;
79 has_work = true;
80 // The condition variable call should be moved after thread_lock.unlock() for performance
81 // reasons but thread sanitizers may report it as a data race if we do.
82 // See https://stackoverflow.com/questions/35775501/c-should-condition-variable-be-notified-under-lock
83 cond_var.notify_one(); // will notify the thread lock that we have work
84 lock.unlock();
85}
86#endif
87
89 dom::parser &_parser,
90 const uint8_t *_buf,
91 size_t _len,
92 size_t _batch_size,
93 stream_format _format
94) noexcept
95 : parser{&_parser},
96 buf{_buf},
97 len{_len},
98 batch_size{_batch_size <= MINIMAL_BATCH_SIZE ? MINIMAL_BATCH_SIZE : _batch_size},
99 format{_format},
100 error{SUCCESS}
101#ifdef SIMDJSON_THREADS_ENABLED
102 , use_thread(_parser.threaded) // we need to make a copy because _parser.threaded can change
103#endif
104{
105#ifdef SIMDJSON_THREADS_ENABLED
106 if(worker.get() == nullptr) {
107 error = MEMALLOC;
108 }
109#endif
110}
111
112simdjson_inline document_stream::document_stream() noexcept
113 : parser{nullptr},
114 buf{nullptr},
115 len{0},
116 batch_size{0},
118 error{UNINITIALIZED}
119#ifdef SIMDJSON_THREADS_ENABLED
120 , use_thread(false)
121#endif
122{
123}
124
125simdjson_inline document_stream::~document_stream() noexcept {
126#ifdef SIMDJSON_THREADS_ENABLED
127 worker.reset();
128#endif
129}
130
131simdjson_inline document_stream::iterator::iterator() noexcept
132 : stream{nullptr}, finished{true} {
133}
134
136 start();
137 // If there are no documents, we're finished.
138 return iterator(this, error == EMPTY);
139}
140
142 return iterator(this, true);
143}
144
145simdjson_inline document_stream::iterator::iterator(document_stream* _stream, bool is_end) noexcept
146 : stream{_stream}, finished{is_end} {
147}
148
150 // Note that in case of error, we do not yet mark
151 // the iterator as "finished": this detection is done
152 // in the operator++ function since it is possible
153 // to call operator++ repeatedly while omitting
154 // calls to operator*.
155 if (stream->error) { return stream->error; }
156 return stream->parser->doc.root();
157}
158
160 // If there is an error, then we want the iterator
161 // to be finished, no matter what. (E.g., we do not
162 // keep generating documents with errors, or go beyond
163 // a document with errors.)
164 //
165 // Users do not have to call "operator*()" when they use operator++,
166 // so we need to end the stream in the operator++ function.
167 //
168 // Note that setting finished = true is essential otherwise
169 // we would enter an infinite loop.
170 if (stream->error) { finished = true; }
171 // Note that stream->error() is guarded against error conditions
172 // (it will immediately return if stream->error casts to false).
173 // In effect, this next function does nothing when (stream->error)
174 // is true (hence the risk of an infinite loop).
175 stream->next();
176 // If that was the last document, we're finished.
177 // It is the only type of error we do not want to appear
178 // in operator*.
179 if (stream->error == EMPTY) { finished = true; }
180 // If we had any other kind of error (not EMPTY) then we want
181 // to pass it along to the operator* and we cannot mark the result
182 // as "finished" just yet.
183 return *this;
184}
185
186simdjson_inline bool document_stream::iterator::operator!=(const document_stream::iterator &other) const noexcept {
187 return finished != other.finished;
188}
189
190inline void document_stream::start() noexcept {
191 if (error) { return; }
192 error = parser->ensure_capacity(batch_size);
193 if (error) { return; }
194 parser->implementation->_number_as_string = parser->number_as_string();
195 // Always run the first stage 1 parse immediately
196 batch_start = 0;
197 error = run_stage1(*parser, batch_start);
198 while(error == EMPTY) {
199 // In exceptional cases, we may start with an empty block
200 batch_start = next_batch_start();
201 if (batch_start >= len) { return; }
202 error = run_stage1(*parser, batch_start);
203 }
204 if (error) { return; }
205#ifdef SIMDJSON_THREADS_ENABLED
206 if (use_thread && next_batch_start() < len) {
207 // Kick off the first thread if needed
208 error = stage1_thread_parser.ensure_capacity(batch_size);
209 if (error) { return; }
210 worker->start_thread();
211 start_stage1_thread();
212 if (error) { return; }
213 }
214#endif // SIMDJSON_THREADS_ENABLED
215 next();
216}
217
218simdjson_inline size_t document_stream::iterator::current_index() const noexcept {
219 return stream->doc_index;
220}
221
222simdjson_inline std::string_view document_stream::iterator::source() const noexcept {
223 const char* start = reinterpret_cast<const char*>(stream->buf) + current_index();
224 if (stream->error) {
225 return std::string_view(start, stream->len - current_index());
226 }
227 bool object_or_array = ((*start == '[') || (*start == '{'));
228 if(object_or_array) {
229 size_t next_doc_index = stream->batch_start + stream->parser->implementation->structural_indexes[stream->parser->implementation->next_structural_index - 1];
230 return std::string_view(start, next_doc_index - current_index() + 1);
231 } else {
232 size_t next_doc_index = stream->batch_start + stream->parser->implementation->structural_indexes[stream->parser->implementation->next_structural_index];
233 size_t svlen = next_doc_index - current_index();
234 // When the scalar is followed by a truncated document, the structural
235 // indexes of that document were dropped and next_doc_index is the end of
236 // the input, so we bound the scalar by scanning the token itself.
237 size_t token_len = 0;
238 if (*start == '"') {
239 token_len = 1;
240 while (token_len < svlen) {
241 char c = start[token_len++];
242 if (c == '\\') {
243 token_len++;
244 } else if (c == '"') {
245 break;
246 }
247 }
248 } else {
249 while (token_len < svlen) {
250 char c = start[token_len];
251 if (std::isspace(static_cast<unsigned char>(c)) || c == ',' || c == '{' || c == '[' || c == '\0' || static_cast<uint8_t>(c) == 0x1E) {
252 break;
253 }
254 token_len++;
255 }
256 }
257 if (token_len > 0 && token_len < svlen) {
258 svlen = token_len;
259 }
260 // Trim trailing whitespace, NUL, and RS (0x1E). In RFC 7464 json_sequence
261 // mode the scanner classifies RS as a scalar character, so an RS-prefixed
262 // scalar document (number/true/false/null/string) has no closing structural
263 // index and the slice runs all the way up to the next document's RS. RS
264 // cannot legally appear in a JSON value at the source level (control
265 // characters in strings must be escaped as \u001E), so stripping it is
266 // safe in every stream_format.
267 while(svlen > 1 && (std::isspace(static_cast<unsigned char>(start[svlen-1])) || start[svlen-1] == '\0' || static_cast<uint8_t>(start[svlen-1]) == 0x1E || (stream->format == stream_format::comma_delimited && start[svlen-1] == ','))) {
268 svlen--;
269 }
270 return std::string_view(start, svlen);
271 }
272}
273
274
275inline void document_stream::next() noexcept {
276 // We always exit at once, once in an error condition.
277 if (error) { return; }
278
279 // Load the next document from the batch
280 doc_index = batch_start + parser->implementation->structural_indexes[parser->implementation->next_structural_index];
281 error = parser->implementation->stage2_next(parser->doc);
282 // If that was the last document in the batch, load another batch (if available)
283 while (error == EMPTY) {
284 batch_start = next_batch_start();
285 if (batch_start >= len) { break; }
286
287#ifdef SIMDJSON_THREADS_ENABLED
288 if(use_thread) {
289 load_from_stage1_thread();
290 } else {
291 error = run_stage1(*parser, batch_start);
292 }
293#else
294 error = run_stage1(*parser, batch_start);
295#endif
296 if (error) { continue; } // If the error was EMPTY, we may want to load another batch.
297 // Run stage 2 on the first document in the batch
298 doc_index = batch_start + parser->implementation->structural_indexes[parser->implementation->next_structural_index];
299 error = parser->implementation->stage2_next(parser->doc);
300 }
301}
302inline size_t document_stream::size_in_bytes() const noexcept {
303 return len;
304}
305
306inline size_t document_stream::truncated_bytes() const noexcept {
307 // Stage 1 returns EMPTY on zero-length input before it writes the index
308 // sentinels read below, so they would still hold a previous stream's values.
309 if (len == 0) { return 0; }
310 if(error == CAPACITY) { return len - batch_start; }
311 return parser->implementation->structural_indexes[parser->implementation->n_structural_indexes] - parser->implementation->structural_indexes[parser->implementation->n_structural_indexes + 1];
312}
313
314inline size_t document_stream::next_batch_start() const noexcept {
315 return batch_start + parser->implementation->structural_indexes[parser->implementation->n_structural_indexes];
316}
317
318inline error_code document_stream::run_stage1(dom::parser &p, size_t _batch_start) noexcept {
319 size_t remaining = len - _batch_start;
320 stage1_mode mode;
321 if (remaining <= batch_size) {
322 // Final batch
323 switch (format) {
325 mode = stage1_mode::json_sequence_final;
326 break;
328 mode = stage1_mode::comma_delimited_final;
329 break;
330 default:
331 mode = stage1_mode::streaming_final;
332 break;
333 }
334 return p.implementation->stage1(&buf[_batch_start], remaining, mode);
335 } else {
336 // Partial batch
337 switch (format) {
339 mode = stage1_mode::json_sequence_partial;
340 break;
342 mode = stage1_mode::comma_delimited_partial;
343 break;
344 default:
345 mode = stage1_mode::streaming_partial;
346 break;
347 }
348 return p.implementation->stage1(&buf[_batch_start], batch_size, mode);
349 }
350}
351
352#ifdef SIMDJSON_THREADS_ENABLED
353
354inline void document_stream::load_from_stage1_thread() noexcept {
355 worker->finish();
356 // Swap to the parser that was loaded up in the thread. Make sure the parser has
357 // enough memory to swap to, as well.
358 std::swap(*parser, stage1_thread_parser);
359 error = stage1_thread_error;
360 if (error) { return; }
361
362 // If there's anything left, start the stage 1 thread!
363 if (next_batch_start() < len) {
364 start_stage1_thread();
365 }
366}
367
368inline void document_stream::start_stage1_thread() noexcept {
369 // we call the thread on a lambda that will update
370 // this->stage1_thread_error
371 // there is only one thread that may write to this value
372 // TODO this is NOT exception-safe.
373 this->stage1_thread_error = UNINITIALIZED; // In case something goes wrong, make sure it's an error
374 size_t _next_batch_start = this->next_batch_start();
375
376 worker->run(this, & this->stage1_thread_parser, _next_batch_start);
377}
378
379#endif // SIMDJSON_THREADS_ENABLED
380
381} // namespace dom
382
383simdjson_inline simdjson_result<dom::document_stream>::simdjson_result() noexcept
384 : simdjson_result_base() {
385}
386simdjson_inline simdjson_result<dom::document_stream>::simdjson_result(error_code error) noexcept
387 : simdjson_result_base(error) {
388}
389simdjson_inline simdjson_result<dom::document_stream>::simdjson_result(dom::document_stream &&value) noexcept
390 : simdjson_result_base(std::forward<dom::document_stream>(value)) {
391}
392
393#if SIMDJSON_EXCEPTIONS
394simdjson_inline dom::document_stream::iterator simdjson_result<dom::document_stream>::begin() noexcept(false) {
395 if (error()) { throw simdjson_error(error()); }
396 return first.begin();
397}
398simdjson_inline dom::document_stream::iterator simdjson_result<dom::document_stream>::end() noexcept(false) {
399 if (error()) { throw simdjson_error(error()); }
400 return first.end();
401}
402#else // SIMDJSON_EXCEPTIONS
403#ifndef SIMDJSON_DISABLE_DEPRECATED_API
404simdjson_inline dom::document_stream::iterator simdjson_result<dom::document_stream>::begin() noexcept {
405 first.error = error();
406 return first.begin();
407}
408simdjson_inline dom::document_stream::iterator simdjson_result<dom::document_stream>::end() noexcept {
409 first.error = error();
410 return first.end();
411}
412#endif // SIMDJSON_DISABLE_DEPRECATED_API
413#endif // SIMDJSON_EXCEPTIONS
414
415} // namespace simdjson
416#endif // SIMDJSON_DOCUMENT_STREAM_INL_H
An iterator through a forward-only stream of documents.
simdjson_inline reference operator*() noexcept
Get the current document (or error).
simdjson_inline bool operator!=(const iterator &other) const noexcept
Check if we're at the end yet.
simdjson_inline iterator() noexcept
Default constructor.
iterator & operator++() noexcept
Advance to the next document (prefix).
A forward-only stream of documents.
size_t size_in_bytes() const noexcept
Returns the input size in bytes.
size_t truncated_bytes() const noexcept
After iterating through the stream, this method returns the number of bytes that were not parsed at t...
simdjson_inline iterator begin() noexcept
Start iterating the documents in the stream.
simdjson_inline iterator end() noexcept
The end of the stream, for iterator comparison purposes.
simdjson_inline document_stream() noexcept
Construct an uninitialized document_stream.
A persistent document parser.
Definition parser.h:30
void number_as_string(bool enabled) noexcept
When enabled, big integers (exceeding uint64 range) are stored as strings in the tape instead of retu...
Definition parser.h:711
The top level simdjson namespace, containing everything the library provides.
Definition base.h:8
stream_format
Stream format for parse_many/iterate_many.
Definition base.h:61
@ comma_delimited
Comma-separated JSON documents (e.g., {...},{...},{...})
@ whitespace_delimited
Whitespace-delimited JSON documents (default, includes NDJSON/JSONL)
@ json_sequence
RFC 7464 JSON text sequences (RS-delimited)
error_code
All possible errors returned by simdjson.
Definition error.h:19
@ CAPACITY
This parser can't support a document that big.
Definition error.h:21
@ EMPTY
no structural element found
Definition error.h:33
@ MEMALLOC
Error allocating memory, most likely out of memory.
Definition error.h:22
@ SUCCESS
No error.
Definition error.h:20
@ UNINITIALIZED
unknown error, or uninitialized document
Definition error.h:32
stage1_mode
This enum is used with the dom_parser_implementation::stage1 function.
The result of a simdjson operation that could fail.
Definition error.h:281
simdjson_inline error_code error() const noexcept
The error.
Definition error-inl.h:168