simdjson 5.0.1
Ridiculously Fast JSON
Loading...
Searching...
No Matches
document_stream-inl.h
1#ifndef SIMDJSON_GENERIC_ONDEMAND_DOCUMENT_STREAM_INL_H
2
3#ifndef SIMDJSON_CONDITIONAL_INCLUDE
4#define SIMDJSON_GENERIC_ONDEMAND_DOCUMENT_STREAM_INL_H
5#include "simdjson/generic/ondemand/base.h"
6#include "simdjson/generic/ondemand/document_stream.h"
7#include "simdjson/generic/ondemand/document-inl.h"
8#include "simdjson/generic/implementation_simdjson_result_base-inl.h"
9#endif // SIMDJSON_CONDITIONAL_INCLUDE
10
11#include <algorithm>
12#include <cstring>
13#include <stdexcept>
14
15namespace simdjson {
16namespace SIMDJSON_IMPLEMENTATION {
17namespace ondemand {
18
19#ifdef SIMDJSON_THREADS_ENABLED
20
21inline void stage1_worker::finish() {
22 // After calling "run" someone would call finish() to wait
23 // for the end of the processing.
24 // This function will wait until either the thread has done
25 // the processing or, else, the destructor has been called.
26 std::unique_lock<std::mutex> lock(locking_mutex);
27 cond_var.wait(lock, [this]{return has_work == false;});
28}
29
30inline stage1_worker::~stage1_worker() {
31 // The thread may never outlive the stage1_worker instance
32 // and will always be stopped/joined before the stage1_worker
33 // instance is gone.
34 stop_thread();
35}
36
37inline void stage1_worker::start_thread() {
38 std::unique_lock<std::mutex> lock(locking_mutex);
39 if(thread.joinable()) {
40 return; // This should never happen but we never want to create more than one thread.
41 }
42 thread = std::thread([this]{
43 while(true) {
44 std::unique_lock<std::mutex> thread_lock(locking_mutex);
45 // We wait for either "run" or "stop_thread" to be called.
46 cond_var.wait(thread_lock, [this]{return has_work || !can_work;});
47 // If, for some reason, the stop_thread() method was called (i.e., the
48 // destructor of stage1_worker is called, then we want to immediately destroy
49 // the thread (and not do any more processing).
50 if(!can_work) {
51 break;
52 }
53 this->owner->stage1_thread_error = this->owner->run_stage1(*this->stage1_thread_parser,
54 this->_next_batch_start);
55 this->has_work = false;
56 // The condition variable call should be moved after thread_lock.unlock() for performance
57 // reasons but thread sanitizers may report it as a data race if we do.
58 // See https://stackoverflow.com/questions/35775501/c-should-condition-variable-be-notified-under-lock
59 cond_var.notify_one(); // will notify "finish"
60 thread_lock.unlock();
61 }
62 }
63 );
64}
65
66
67inline void stage1_worker::stop_thread() {
68 std::unique_lock<std::mutex> lock(locking_mutex);
69 // We have to make sure that all locks can be released.
70 can_work = false;
71 has_work = false;
72 cond_var.notify_all();
73 lock.unlock();
74 if(thread.joinable()) {
75 thread.join();
76 }
77}
78
79inline void stage1_worker::run(document_stream * ds, parser * stage1, size_t next_batch_start) {
80 std::unique_lock<std::mutex> lock(locking_mutex);
81 owner = ds;
82 _next_batch_start = next_batch_start;
83 stage1_thread_parser = stage1;
84 has_work = true;
85 // The condition variable call should be moved after thread_lock.unlock() for performance
86 // reasons but thread sanitizers may report it as a data race if we do.
87 // See https://stackoverflow.com/questions/35775501/c-should-condition-variable-be-notified-under-lock
88 cond_var.notify_one(); // will notify the thread lock that we have work
89 lock.unlock();
90}
91
92#endif // SIMDJSON_THREADS_ENABLED
93
95 ondemand::parser &_parser,
96 const uint8_t *_buf,
97 size_t _len,
98 size_t _batch_size,
99 bool _allow_comma_separated,
100 stream_format _format
101) noexcept
102 : parser{&_parser},
103 buf{_buf},
104 len{_len},
105 batch_size{_batch_size <= MINIMAL_BATCH_SIZE ? MINIMAL_BATCH_SIZE : _batch_size},
106 allow_comma_separated{_allow_comma_separated},
107 format{_format},
108 error{SUCCESS}
109 #ifdef SIMDJSON_THREADS_ENABLED
110 , use_thread(_parser.threaded) // we need to make a copy because _parser.threaded can change
111 #endif
112{
113}
114
115simdjson_inline document_stream::document_stream() noexcept
116 : parser{nullptr},
117 buf{nullptr},
118 len{0},
119 batch_size{0},
120 allow_comma_separated{false},
122 error{UNINITIALIZED}
123 #ifdef SIMDJSON_THREADS_ENABLED
124 , use_thread(false)
125 #endif
126{
127}
128
129simdjson_inline document_stream::~document_stream() noexcept
130{
131 #ifdef SIMDJSON_THREADS_ENABLED
132 worker.reset();
133 #endif
134}
135
136inline size_t document_stream::size_in_bytes() const noexcept {
137 return len;
138}
139
140inline size_t document_stream::truncated_bytes() const noexcept {
141 if(error == CAPACITY) { return len - batch_start; }
142 return parser->implementation->structural_indexes[parser->implementation->n_structural_indexes] - parser->implementation->structural_indexes[parser->implementation->n_structural_indexes + 1];
143}
144
145simdjson_inline document_stream::iterator::iterator() noexcept
146 : stream{nullptr}, finished{true} {
147}
148
149simdjson_inline document_stream::iterator::iterator(document_stream* _stream, bool is_end) noexcept
150 : stream{_stream}, finished{is_end} {
151}
152
156
158 // If there is an error, then we want the iterator
159 // to be finished, no matter what. (E.g., we do not
160 // keep generating documents with errors, or go beyond
161 // a document with errors.)
162 //
163 // Users do not have to call "operator*()" when they use operator++,
164 // so we need to end the stream in the operator++ function.
165 //
166 // Note that setting finished = true is essential otherwise
167 // we would enter an infinite loop.
168 if (stream->error) { finished = true; }
169 // Note that stream->error() is guarded against error conditions
170 // (it will immediately return if stream->error casts to false).
171 // In effect, this next function does nothing when (stream->error)
172 // is true (hence the risk of an infinite loop).
173 stream->next();
174 // If that was the last document, we're finished.
175 // It is the only type of error we do not want to appear
176 // in operator*.
177 if (stream->error == EMPTY) { finished = true; }
178 // If we had any other kind of error (not EMPTY) then we want
179 // to pass it along to the operator* and we cannot mark the result
180 // as "finished" just yet.
181 return *this;
182}
183
184simdjson_inline bool document_stream::iterator::at_end() const noexcept {
185 return finished;
186}
187
188
189simdjson_inline bool document_stream::iterator::operator!=(const document_stream::iterator &other) const noexcept {
190 return finished != other.finished;
191}
192
193simdjson_inline bool document_stream::iterator::operator==(const document_stream::iterator &other) const noexcept {
194 return finished == other.finished;
195}
196
198 start();
199 // If there are no documents, we're finished.
200 return iterator(this, error == EMPTY);
201}
202
204 return iterator(this, true);
205}
206
207inline void document_stream::start() noexcept {
208 if (error) { return; }
209 error = parser->allocate(batch_size);
210 if (error) { return; }
211 // Always run the first stage 1 parse immediately
212 batch_start = 0;
213 error = run_stage1(*parser, batch_start);
214 while(error == EMPTY) {
215 // In exceptional cases, we may start with an empty block
216 batch_start = next_batch_start();
217 if (batch_start >= len) { return; }
218 error = run_stage1(*parser, batch_start);
219 }
220 if (error) { return; }
221 // For json_sequence mode, structural_indexes[0] points to the actual JSON value
222 // after the RS delimiter and any following whitespace. For regular mode, it is
223 // the offset from batch_start to the first document in the batch.
224 doc_index = batch_start + parser->implementation->structural_indexes[0];
225 doc = document(json_iterator(&buf[batch_start], parser));
226 doc.iter._streaming = true;
227
228 #ifdef SIMDJSON_THREADS_ENABLED
229 if (use_thread && next_batch_start() < len) {
230 // Kick off the first thread on next batch if needed
231 if (worker.get() == nullptr) {
232 worker.reset(new(std::nothrow) stage1_worker());
233 if (worker.get() == nullptr) { error = MEMALLOC; return; }
234 }
235 error = stage1_thread_parser.allocate(batch_size);
236 if (error) { return; }
237 worker->start_thread();
238 start_stage1_thread();
239 if (error) { return; }
240 }
241 #endif // SIMDJSON_THREADS_ENABLED
242}
243
244inline void document_stream::next() noexcept {
245 // We always enter at once once in an error condition.
246 if (error) { return; }
247 next_document();
248 if (error) { return; }
249 auto cur_struct_index = doc.iter._root - parser->implementation->structural_indexes.get();
250 doc_index = batch_start + parser->implementation->structural_indexes[cur_struct_index];
251
252 // Check if at end of structural indexes (i.e. at end of batch)
253 if(cur_struct_index >= static_cast<int64_t>(parser->implementation->n_structural_indexes)) {
254 error = EMPTY;
255 // Load another batch (if available)
256 while (error == EMPTY) {
257 batch_start = next_batch_start();
258 if (batch_start >= len) { break; }
259 #ifdef SIMDJSON_THREADS_ENABLED
260 if(use_thread) {
261 load_from_stage1_thread();
262 } else {
263 error = run_stage1(*parser, batch_start);
264 }
265 #else
266 error = run_stage1(*parser, batch_start);
267 #endif
302 doc.iter = json_iterator(&buf[batch_start], parser);
303 doc.iter._streaming = true;
308 if (error) { continue; } // If the error was EMPTY, we may want to load another batch.
309 doc_index = batch_start + parser->implementation->structural_indexes[0];
310 }
311 }
312}
313
314simdjson_inline uint8_t document_stream::document_delimiter() const noexcept {
315 switch (format) {
316 case stream_format::newline_delimited: return '\n';
317 case stream_format::json_sequence: return 0x1E;
318 default: return 0;
319 }
320}
321
322simdjson_inline bool document_stream::skip_to_delimiter(uint8_t delimiter) noexcept {
323 const uint8_t *const base = &buf[batch_start];
324 const token_position pos = doc.iter.position();
325 const token_position end = doc.iter.end_position();
326 if (pos >= end) { return false; }
327 const size_t here = size_t(doc.iter.token.peek(pos) - base);
328 const size_t batch_len =
329 (len - batch_start < batch_size) ? len - batch_start : batch_size;
330 if (here >= batch_len) { return false; }
331 const uint8_t *const found = static_cast<const uint8_t *>(
332 std::memchr(base + here, delimiter, batch_len - here));
333 if (found == nullptr) { return false; }
334
335 const uint32_t boundary = uint32_t(found - base);
336 // The answer is near `pos`: the delimiter ends the current document, while
337 // `end` spans the whole batch. Gallop first so the cost follows the distance
338 // rather than the size of the batch.
339 token_position lo = pos;
340 size_t hop = 1;
341 while (lo + hop < end && lo[hop] < boundary) { lo += hop; hop <<= 1; }
342 token_position hi = (lo + hop < end) ? lo + hop : end;
343 while (lo < hi) {
344 const token_position mid = lo + ((hi - lo) >> 1);
345 if (*mid < boundary) { lo = mid + 1; } else { hi = mid; }
346 }
347 doc.iter.token.set_position(lo);
348 return true;
349}
350
351inline void document_stream::next_document() noexcept {
352 // A delimiter that cannot occur inside a document tells us where the current
353 // one ends, so we can jump there instead of walking every structural. Only
354 // valid while the iterator is still inside the document: a consumed document
355 // already sits on the next one's first token, and skip_child() returns at
356 // once for it.
357 //
358 // The jump does not structure-validate the unread remainder of the current
359 // document: under newline_delimited / json_sequence the next delimiter is
360 // assumed to be the true document boundary. Callers that leave depth() > 0
361 // while violating that contract (e.g. pretty multi-line JSON under
362 // newline_delimited) can mis-align following documents; use
363 // whitespace_delimited if unsure.
364 const uint8_t delimiter = document_delimiter();
365 if (delimiter != 0 && !error && doc.iter.depth() > 0 &&
366 skip_to_delimiter(delimiter)) {
367 doc.iter._depth = 1;
368 doc.iter._string_buf_loc = parser->string_buf.get();
369 doc.iter._root = doc.iter.position();
370 return;
371 }
372 // Go to next place where depth=0 (document depth)
373 error = doc.iter.skip_child(0);
374 if (error) { return; }
375 // Always set depth=1 at the start of document
376 doc.iter._depth = 1;
377 // consume comma if comma separated is allowed
378 if (allow_comma_separated) {
379 error_code ignored = doc.iter.consume_character(',');
380 static_cast<void>(ignored); // ignored on purpose
381 }
382 // Resets the string buffer at the beginning, thus invalidating the strings.
383 doc.iter._string_buf_loc = parser->string_buf.get();
384 doc.iter._root = doc.iter.position();
385}
386
387inline size_t document_stream::next_batch_start() const noexcept {
388 return batch_start + parser->implementation->structural_indexes[parser->implementation->n_structural_indexes];
389}
390
391inline error_code document_stream::run_stage1(ondemand::parser &p, size_t _batch_start) noexcept {
392 // This code only updates the structural index in the parser, it does not update any json_iterator
393 // instance.
394 size_t remaining = len - _batch_start;
395 stage1_mode mode;
396 if (remaining <= batch_size) {
397 // Final batch
398 switch (format) {
400 mode = stage1_mode::json_sequence_final;
401 break;
403 mode = stage1_mode::comma_delimited_final;
404 break;
405 default:
406 mode = stage1_mode::streaming_final;
407 break;
408 }
409 return p.implementation->stage1(&buf[_batch_start], remaining, mode);
410 } else {
411 // Partial batch
412 switch (format) {
414 mode = stage1_mode::json_sequence_partial;
415 break;
417 mode = stage1_mode::comma_delimited_partial;
418 break;
419 default:
420 mode = stage1_mode::streaming_partial;
421 break;
422 }
423 return p.implementation->stage1(&buf[_batch_start], batch_size, mode);
424 }
425}
426
427simdjson_inline size_t document_stream::iterator::current_index() const noexcept {
428 return stream->doc_index;
429}
430
431simdjson_inline std::string_view document_stream::iterator::source() const noexcept {
432 auto depth = stream->doc.iter.depth();
433 auto cur_struct_index = stream->doc.iter._root - stream->parser->implementation->structural_indexes.get();
434
435 // If at root, process the first token to determine if scalar value
436 if (stream->doc.iter.at_root()) {
437 switch (stream->buf[stream->batch_start + stream->parser->implementation->structural_indexes[cur_struct_index]]) {
438 case '{': case '[': // Depth=1 already at start of document
439 break;
440 case '}': case ']':
441 depth--;
442 break;
443 default: // Scalar value document
444 // This returns a string spanning from start of value to the beginning of the next document (excluded)
445 {
446 auto next_index = stream->batch_start + stream->parser->implementation->structural_indexes[++cur_struct_index];
447 // normally the length would be next_index - current_index() - 1, except for the last document
448 size_t svlen = next_index - current_index();
449 const char *start = reinterpret_cast<const char*>(stream->buf) + current_index();
450 // Trim trailing whitespace, NUL, and RS (0x1E). In RFC 7464
451 // json_sequence mode the scanner classifies RS as a scalar
452 // character, so an RS-prefixed scalar document (number / true /
453 // false / null / string) has no closing structural index and the
454 // slice runs all the way up to the next document's RS. RS cannot
455 // legally appear in a JSON value at the source level (control
456 // characters in strings must be escaped as \u001E), so stripping
457 // it is safe in every stream_format.
458 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] == ','))) {
459 svlen--;
460 }
461 return std::string_view(start, svlen);
462 }
463 }
464 cur_struct_index++;
465 }
466
467 while (cur_struct_index <= static_cast<int64_t>(stream->parser->implementation->n_structural_indexes)) {
468 switch (stream->buf[stream->batch_start + stream->parser->implementation->structural_indexes[cur_struct_index]]) {
469 case '{': case '[':
470 depth++;
471 break;
472 case '}': case ']':
473 depth--;
474 break;
475 }
476 if (depth == 0) { break; }
477 cur_struct_index++;
478 }
479
480 return std::string_view(reinterpret_cast<const char*>(stream->buf) + current_index(), stream->parser->implementation->structural_indexes[cur_struct_index] - current_index() + stream->batch_start + 1);;
481}
482
484 return stream->error;
485}
486
487#ifdef SIMDJSON_THREADS_ENABLED
488
489inline void document_stream::load_from_stage1_thread() noexcept {
490 worker->finish();
491 // Swap to the parser that was loaded up in the thread. Make sure the parser has
492 // enough memory to swap to, as well.
493 std::swap(stage1_thread_parser,*parser);
494 error = stage1_thread_error;
495 if (error) { return; }
496
497 // If there's anything left, start the stage 1 thread!
498 if (next_batch_start() < len) {
499 start_stage1_thread();
500 }
501}
502
503inline void document_stream::start_stage1_thread() noexcept {
504 // we call the thread on a lambda that will update
505 // this->stage1_thread_error
506 // there is only one thread that may write to this value
507 // TODO this is NOT exception-safe.
508 this->stage1_thread_error = UNINITIALIZED; // In case something goes wrong, make sure it's an error
509 size_t _next_batch_start = this->next_batch_start();
510
511 worker->run(this, & this->stage1_thread_parser, _next_batch_start);
512}
513
514#endif // SIMDJSON_THREADS_ENABLED
515
516} // namespace ondemand
517} // namespace SIMDJSON_IMPLEMENTATION
518} // namespace simdjson
519
520namespace simdjson {
521
522simdjson_inline simdjson_result<SIMDJSON_IMPLEMENTATION::ondemand::document_stream>::simdjson_result(
523 error_code error
524) noexcept :
525 implementation_simdjson_result_base<SIMDJSON_IMPLEMENTATION::ondemand::document_stream>(error)
526{
527}
528simdjson_inline simdjson_result<SIMDJSON_IMPLEMENTATION::ondemand::document_stream>::simdjson_result(
529 SIMDJSON_IMPLEMENTATION::ondemand::document_stream &&value
530) noexcept :
531 implementation_simdjson_result_base<SIMDJSON_IMPLEMENTATION::ondemand::document_stream>(
532 std::forward<SIMDJSON_IMPLEMENTATION::ondemand::document_stream>(value)
533 )
534{
535}
536
537}
538
539#endif // SIMDJSON_GENERIC_ONDEMAND_DOCUMENT_STREAM_INL_H
iterator & operator++() noexcept
Advance to the next document (prefix).
simdjson_inline bool operator!=(const iterator &other) const noexcept
Check if we're at the end yet.
error_code error() const noexcept
Returns error of the stream (if any).
simdjson_inline reference operator*() noexcept
Get the current document (or error).
bool at_end() const noexcept
Returns whether the iterator is at the end.
simdjson_inline iterator end() noexcept
The end of the stream, for iterator comparison purposes.
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 document_stream() noexcept
Construct an uninitialized document_stream.
json_iterator iter
Current position in the document.
Definition document.h:934
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)
@ newline_delimited
NDJSON/JSON Lines where each document occupies exactly one line: documents are separated by line feed...
@ 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