1#ifndef SIMDJSON_DOCUMENT_STREAM_INL_H
2#define SIMDJSON_DOCUMENT_STREAM_INL_H
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"
14#ifdef SIMDJSON_THREADS_ENABLED
16inline void stage1_worker::finish() {
21 std::unique_lock<std::mutex> lock(locking_mutex);
22 cond_var.wait(lock, [
this]{
return has_work ==
false;});
25inline stage1_worker::~stage1_worker() {
32inline void stage1_worker::start_thread() {
33 std::unique_lock<std::mutex> lock(locking_mutex);
34 if(thread.joinable()) {
37 thread = std::thread([
this]{
39 std::unique_lock<std::mutex> thread_lock(locking_mutex);
41 cond_var.wait(thread_lock, [
this]{
return has_work || !can_work;});
48 this->owner->stage1_thread_error = this->owner->run_stage1(*this->stage1_thread_parser,
49 this->_next_batch_start);
50 this->has_work =
false;
54 cond_var.notify_one();
62inline void stage1_worker::stop_thread() {
63 std::unique_lock<std::mutex> lock(locking_mutex);
67 cond_var.notify_all();
69 if(thread.joinable()) {
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);
77 _next_batch_start = next_batch_start;
78 stage1_thread_parser = stage1;
83 cond_var.notify_one();
98 batch_size{_batch_size <= MINIMAL_BATCH_SIZE ? MINIMAL_BATCH_SIZE : _batch_size},
101#ifdef SIMDJSON_THREADS_ENABLED
102 , use_thread(_parser.threaded)
105#ifdef SIMDJSON_THREADS_ENABLED
106 if(worker.get() ==
nullptr) {
119#ifdef SIMDJSON_THREADS_ENABLED
125simdjson_inline document_stream::~document_stream() noexcept {
126#ifdef SIMDJSON_THREADS_ENABLED
132 : stream{
nullptr}, finished{
true} {
146 : stream{_stream}, finished{is_end} {
155 if (stream->error) {
return stream->
error; }
156 return stream->parser->doc.root();
170 if (stream->error) { finished =
true; }
179 if (stream->error ==
EMPTY) { finished =
true; }
187 return finished != other.finished;
190inline void document_stream::start() noexcept {
191 if (error) {
return; }
192 error =
parser->ensure_capacity(batch_size);
193 if (error) {
return; }
197 error = run_stage1(*parser, batch_start);
198 while(error ==
EMPTY) {
200 batch_start = next_batch_start();
201 if (batch_start >= len) {
return; }
202 error = run_stage1(*parser, batch_start);
204 if (error) {
return; }
205#ifdef SIMDJSON_THREADS_ENABLED
206 if (use_thread && next_batch_start() < len) {
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; }
218simdjson_inline
size_t document_stream::iterator::current_index() const noexcept {
219 return stream->doc_index;
222simdjson_inline std::string_view document_stream::iterator::source() const noexcept {
223 const char* start =
reinterpret_cast<const char*
>(stream->buf) + current_index();
225 return std::string_view(start, stream->len - current_index());
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);
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();
237 size_t token_len = 0;
240 while (token_len < svlen) {
241 char c = start[token_len++];
244 }
else if (c ==
'"') {
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) {
257 if (token_len > 0 && token_len < svlen) {
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] ==
','))) {
270 return std::string_view(start, svlen);
275inline void document_stream::next() noexcept {
277 if (error) {
return; }
280 doc_index = batch_start + parser->implementation->structural_indexes[parser->implementation->next_structural_index];
281 error = parser->implementation->stage2_next(parser->doc);
283 while (error ==
EMPTY) {
284 batch_start = next_batch_start();
285 if (batch_start >= len) {
break; }
287#ifdef SIMDJSON_THREADS_ENABLED
289 load_from_stage1_thread();
291 error = run_stage1(*parser, batch_start);
294 error = run_stage1(*parser, batch_start);
296 if (error) {
continue; }
298 doc_index = batch_start + parser->implementation->structural_indexes[parser->implementation->next_structural_index];
299 error = parser->implementation->stage2_next(parser->doc);
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];
314inline size_t document_stream::next_batch_start() const noexcept {
315 return batch_start +
parser->implementation->structural_indexes[
parser->implementation->n_structural_indexes];
319 size_t remaining = len - _batch_start;
321 if (remaining <= batch_size) {
325 mode = stage1_mode::json_sequence_final;
328 mode = stage1_mode::comma_delimited_final;
331 mode = stage1_mode::streaming_final;
334 return p.implementation->stage1(&buf[_batch_start], remaining, mode);
339 mode = stage1_mode::json_sequence_partial;
342 mode = stage1_mode::comma_delimited_partial;
345 mode = stage1_mode::streaming_partial;
348 return p.implementation->stage1(&buf[_batch_start], batch_size, mode);
352#ifdef SIMDJSON_THREADS_ENABLED
354inline void document_stream::load_from_stage1_thread() noexcept {
358 std::swap(*parser, stage1_thread_parser);
359 error = stage1_thread_error;
360 if (error) {
return; }
363 if (next_batch_start() < len) {
364 start_stage1_thread();
368inline void document_stream::start_stage1_thread() noexcept {
374 size_t _next_batch_start = this->next_batch_start();
376 worker->run(
this, & this->stage1_thread_parser, _next_batch_start);
383simdjson_inline simdjson_result<dom::document_stream>::simdjson_result() noexcept
384 : simdjson_result_base() {
386simdjson_inline simdjson_result<dom::document_stream>::simdjson_result(
error_code error) noexcept
387 : simdjson_result_base(error) {
389simdjson_inline simdjson_result<dom::document_stream>::simdjson_result(dom::document_stream &&value) noexcept
390 : simdjson_result_base(std::forward<dom::document_stream>(value)) {
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();
398simdjson_inline dom::document_stream::iterator simdjson_result<dom::document_stream>::end() noexcept(false) {
399 if (error()) {
throw simdjson_error(error()); }
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();
408simdjson_inline dom::document_stream::iterator simdjson_result<dom::document_stream>::end() noexcept {
409 first.error = error();
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.
void number_as_string(bool enabled) noexcept
When enabled, big integers (exceeding uint64 range) are stored as strings in the tape instead of retu...
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)
@ 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.
simdjson_inline error_code error() const noexcept
The error.