1#ifndef SIMDJSON_GENERIC_ONDEMAND_DOCUMENT_STREAM_INL_H
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"
16namespace SIMDJSON_IMPLEMENTATION {
19#ifdef SIMDJSON_THREADS_ENABLED
21inline void stage1_worker::finish() {
26 std::unique_lock<std::mutex> lock(locking_mutex);
27 cond_var.wait(lock, [
this]{
return has_work ==
false;});
30inline stage1_worker::~stage1_worker() {
37inline void stage1_worker::start_thread() {
38 std::unique_lock<std::mutex> lock(locking_mutex);
39 if(thread.joinable()) {
42 thread = std::thread([
this]{
44 std::unique_lock<std::mutex> thread_lock(locking_mutex);
46 cond_var.wait(thread_lock, [
this]{
return has_work || !can_work;});
53 this->owner->stage1_thread_error = this->owner->run_stage1(*this->stage1_thread_parser,
54 this->_next_batch_start);
55 this->has_work =
false;
59 cond_var.notify_one();
67inline void stage1_worker::stop_thread() {
68 std::unique_lock<std::mutex> lock(locking_mutex);
72 cond_var.notify_all();
74 if(thread.joinable()) {
79inline void stage1_worker::run(document_stream * ds, parser * stage1,
size_t next_batch_start) {
80 std::unique_lock<std::mutex> lock(locking_mutex);
82 _next_batch_start = next_batch_start;
83 stage1_thread_parser = stage1;
88 cond_var.notify_one();
95 ondemand::parser &_parser,
99 bool _allow_comma_separated,
105 batch_size{_batch_size <= MINIMAL_BATCH_SIZE ? MINIMAL_BATCH_SIZE : _batch_size},
106 allow_comma_separated{_allow_comma_separated},
109 #ifdef SIMDJSON_THREADS_ENABLED
110 , use_thread(_parser.threaded)
120 allow_comma_separated{
false},
123 #ifdef SIMDJSON_THREADS_ENABLED
129simdjson_inline document_stream::~document_stream() noexcept
131 #ifdef SIMDJSON_THREADS_ENABLED
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];
146 : stream{
nullptr}, finished{
true} {
150 : stream{_stream}, finished{is_end} {
168 if (stream->error) { finished =
true; }
177 if (stream->error ==
EMPTY) { finished =
true; }
190 return finished != other.finished;
194 return finished == other.finished;
207inline void document_stream::start() noexcept {
208 if (error) {
return; }
209 error = parser->allocate(batch_size);
210 if (error) {
return; }
213 error = run_stage1(*parser, batch_start);
214 while(error ==
EMPTY) {
216 batch_start = next_batch_start();
217 if (batch_start >= len) {
return; }
218 error = run_stage1(*parser, batch_start);
220 if (error) {
return; }
224 doc_index = batch_start + parser->implementation->structural_indexes[0];
225 doc = document(json_iterator(&buf[batch_start], parser));
226 doc.
iter._streaming =
true;
228 #ifdef SIMDJSON_THREADS_ENABLED
229 if (use_thread && next_batch_start() < len) {
231 if (worker.get() ==
nullptr) {
232 worker.reset(
new(std::nothrow) stage1_worker());
233 if (worker.get() ==
nullptr) { error =
MEMALLOC;
return; }
235 error = stage1_thread_parser.allocate(batch_size);
236 if (error) {
return; }
237 worker->start_thread();
238 start_stage1_thread();
239 if (error) {
return; }
244inline void document_stream::next() noexcept {
246 if (error) {
return; }
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];
253 if(cur_struct_index >=
static_cast<int64_t
>(parser->implementation->n_structural_indexes)) {
256 while (error ==
EMPTY) {
257 batch_start = next_batch_start();
258 if (batch_start >= len) {
break; }
259 #ifdef SIMDJSON_THREADS_ENABLED
261 load_from_stage1_thread();
263 error = run_stage1(*parser, batch_start);
266 error = run_stage1(*parser, batch_start);
302 doc.
iter = json_iterator(&buf[batch_start], parser);
303 doc.
iter._streaming =
true;
308 if (error) {
continue; }
309 doc_index = batch_start + parser->implementation->structural_indexes[0];
314simdjson_inline uint8_t document_stream::document_delimiter() const noexcept {
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; }
335 const uint32_t boundary = uint32_t(found - base);
339 token_position lo = pos;
341 while (lo + hop <
end && lo[hop] < boundary) { lo += hop; hop <<= 1; }
342 token_position hi = (lo + hop <
end) ? lo + hop :
end;
344 const token_position mid = lo + ((hi - lo) >> 1);
345 if (*mid < boundary) { lo = mid + 1; }
else { hi = mid; }
347 doc.
iter.token.set_position(lo);
351inline void document_stream::next_document() noexcept {
364 const uint8_t delimiter = document_delimiter();
365 if (delimiter != 0 && !error && doc.
iter.depth() > 0 &&
366 skip_to_delimiter(delimiter)) {
368 doc.
iter._string_buf_loc = parser->string_buf.get();
369 doc.
iter._root = doc.
iter.position();
373 error = doc.
iter.skip_child(0);
374 if (error) {
return; }
378 if (allow_comma_separated) {
380 static_cast<void>(ignored);
383 doc.
iter._string_buf_loc = parser->string_buf.get();
384 doc.
iter._root = doc.
iter.position();
387inline size_t document_stream::next_batch_start() const noexcept {
388 return batch_start + parser->implementation->structural_indexes[parser->implementation->n_structural_indexes];
391inline error_code document_stream::run_stage1(ondemand::parser &p,
size_t _batch_start)
noexcept {
394 size_t remaining = len - _batch_start;
396 if (remaining <= batch_size) {
400 mode = stage1_mode::json_sequence_final;
403 mode = stage1_mode::comma_delimited_final;
406 mode = stage1_mode::streaming_final;
409 return p.implementation->stage1(&buf[_batch_start], remaining, mode);
414 mode = stage1_mode::json_sequence_partial;
417 mode = stage1_mode::comma_delimited_partial;
420 mode = stage1_mode::streaming_partial;
423 return p.implementation->stage1(&buf[_batch_start], batch_size, mode);
427simdjson_inline
size_t document_stream::iterator::current_index() const noexcept {
428 return stream->doc_index;
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();
436 if (stream->doc.iter.at_root()) {
437 switch (stream->buf[stream->batch_start + stream->parser->implementation->structural_indexes[cur_struct_index]]) {
446 auto next_index = stream->batch_start + stream->parser->implementation->structural_indexes[++cur_struct_index];
448 size_t svlen = next_index - current_index();
449 const char *start =
reinterpret_cast<const char*
>(stream->buf) + current_index();
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] ==
','))) {
461 return std::string_view(start, svlen);
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]]) {
476 if (depth == 0) {
break; }
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);;
484 return stream->error;
487#ifdef SIMDJSON_THREADS_ENABLED
489inline void document_stream::load_from_stage1_thread() noexcept {
493 std::swap(stage1_thread_parser,*
parser);
494 error = stage1_thread_error;
495 if (error) {
return; }
498 if (next_batch_start() < len) {
499 start_stage1_thread();
503inline void document_stream::start_stage1_thread() noexcept {
509 size_t _next_batch_start = this->next_batch_start();
511 worker->run(
this, & this->stage1_thread_parser, _next_batch_start);
522simdjson_inline simdjson_result<SIMDJSON_IMPLEMENTATION::ondemand::document_stream>::simdjson_result(
525 implementation_simdjson_result_base<SIMDJSON_IMPLEMENTATION::ondemand::document_stream>(error)
528simdjson_inline simdjson_result<SIMDJSON_IMPLEMENTATION::ondemand::document_stream>::simdjson_result(
529 SIMDJSON_IMPLEMENTATION::ondemand::document_stream &&value
531 implementation_simdjson_result_base<SIMDJSON_IMPLEMENTATION::ondemand::document_stream>(
532 std::forward<SIMDJSON_IMPLEMENTATION::ondemand::document_stream>(value)
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.
simdjson_inline iterator() noexcept
Default constructor.
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.
A JSON fragment iterator.
The top level simdjson namespace, containing everything the library provides.
stream_format
Stream format for parse_many/iterate_many.
@ 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.
@ CAPACITY
This parser can't support a document that big.
@ EMPTY
no structural element found
@ MEMALLOC
Error allocating memory, most likely out of memory.
@ UNINITIALIZED
unknown error, or uninitialized document
stage1_mode
This enum is used with the dom_parser_implementation::stage1 function.
The result of a simdjson operation that could fail.