use crate::ParserConfig;
use crate::error::Result;
use rayon::prelude::*;
pub fn parse<T>(input: &str) -> Result<Vec<T>>
where
T: serde_core::de::DeserializeOwned + Send + 'static,
{
let config = ParserConfig::default();
parse_with_config(input, &config)
}
pub fn parse_with_config<T>(input: &str, config: &ParserConfig) -> Result<Vec<T>>
where
T: serde_core::de::DeserializeOwned + Send + 'static,
{
const SEQUENTIAL_DOCUMENTS: usize = 4;
crate::doc_boundary::validate_document_budget(input, config.max_documents)?;
let mut chunks = crate::doc_boundary::DocumentStream::new(input, config.max_documents);
let mut prefix = Vec::with_capacity(SEQUENTIAL_DOCUMENTS);
while prefix.len() < SEQUENTIAL_DOCUMENTS {
match chunks.next() {
Some(Ok(chunk)) => prefix.push(chunk),
Some(Err(error)) => return Err(error),
None => break,
}
}
if prefix.is_empty() && !input.is_empty() {
if config.max_documents == 0 {
return Err(crate::Error::Budget(crate::BudgetBreach::MaxDocuments {
limit: 0,
observed: 1,
}));
}
prefix.push(input);
}
if prefix.len() < SEQUENTIAL_DOCUMENTS {
return prefix
.iter()
.map(|chunk| crate::from_str_with_config::<T>(chunk, config))
.collect();
}
let mut parsed = prefix
.into_iter()
.map(Ok)
.chain(chunks)
.enumerate()
.par_bridge()
.map(|(index, chunk)| {
(
index,
chunk.and_then(|chunk| crate::from_str_with_config::<T>(chunk, config)),
)
})
.collect::<Vec<_>>();
parsed.sort_unstable_by_key(|(index, _)| *index);
parsed.into_iter().map(|(_, result)| result).collect()
}
pub fn values(input: &str) -> Result<Vec<crate::Value>> {
parse::<crate::Value>(input)
}
pub fn values_with_config(input: &str, config: &ParserConfig) -> Result<Vec<crate::Value>> {
parse_with_config::<crate::Value>(input, config)
}
#[must_use]
pub fn split(input: &str) -> Vec<&str> {
let mut docs = crate::doc_boundary::split_documents(input);
if docs.is_empty() && !input.is_empty() {
docs.push(input);
}
docs
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn split_separates_three_records() {
let yaml = "---\nid: 1\n---\nid: 2\n---\nid: 3\n";
let docs = split(yaml);
assert_eq!(docs.len(), 3);
assert!(docs[0].contains("id: 1"));
assert!(docs[1].contains("id: 2"));
assert!(docs[2].contains("id: 3"));
}
#[test]
fn split_handles_no_separators() {
let yaml = "single: doc\n";
let docs = split(yaml);
assert_eq!(docs.len(), 1);
assert_eq!(docs[0], yaml);
}
#[test]
fn split_handles_empty_input() {
assert!(split("").is_empty());
}
#[test]
fn split_handles_implicit_first_doc() {
let yaml = "name: a\n---\nname: b\n";
let docs = split(yaml);
assert_eq!(docs.len(), 2);
assert!(docs[0].contains("name: a"));
assert!(docs[1].contains("name: b"));
}
#[test]
fn split_ignores_dashes_mid_line() {
let yaml = "key: value---suffix\n";
let docs = split(yaml);
assert_eq!(docs.len(), 1);
assert!(docs[0].contains("value---suffix"));
}
#[test]
fn split_requires_post_marker_whitespace() {
let yaml = "key: a\n---foo\nkey: b\n";
let docs = split(yaml);
assert_eq!(docs.len(), 1, "got: {docs:?}");
}
#[cfg_attr(
miri,
ignore = "rayon/crossbeam-epoch uses int-to-ptr casts unsupported under -Zmiri-strict-provenance"
)]
#[test]
fn parse_round_trips_typed_records() {
#[derive(Debug, serde::Deserialize, PartialEq)]
struct Record {
id: u32,
}
let yaml = "---\nid: 1\n---\nid: 2\n---\nid: 3\n";
let records: Vec<Record> = parse(yaml).unwrap();
assert_eq!(
records,
vec![Record { id: 1 }, Record { id: 2 }, Record { id: 3 }]
);
}
#[cfg_attr(
miri,
ignore = "rayon/crossbeam-epoch uses int-to-ptr casts unsupported under -Zmiri-strict-provenance"
)]
#[test]
fn values_yields_value_per_document() {
let yaml = "---\na: 1\n---\nb: 2\n";
let docs = values(yaml).unwrap();
assert_eq!(docs.len(), 2);
assert_eq!(docs[0]["a"].as_i64(), Some(1));
assert_eq!(docs[1]["b"].as_i64(), Some(2));
}
#[test]
fn parse_with_config_enforces_document_count_before_scheduling() {
let config = ParserConfig::default().max_documents(1);
let error =
parse_with_config::<crate::Value>("---\na: 1\n---\nb: 2\n", &config).unwrap_err();
assert!(matches!(
error,
crate::Error::Budget(crate::BudgetBreach::MaxDocuments { limit: 1, .. })
));
}
#[test]
fn document_overflow_is_rejected_before_deserialization() {
use core::sync::atomic::{AtomicUsize, Ordering};
static DESERIALIZATIONS: AtomicUsize = AtomicUsize::new(0);
#[derive(Debug)]
struct CountingValue;
impl<'de> serde_core::Deserialize<'de> for CountingValue {
fn deserialize<D>(deserializer: D) -> core::result::Result<Self, D::Error>
where
D: serde_core::Deserializer<'de>,
{
let _ = <crate::Value as serde_core::Deserialize>::deserialize(deserializer)?;
let _ = DESERIALIZATIONS.fetch_add(1, Ordering::Relaxed);
Ok(Self)
}
}
DESERIALIZATIONS.store(0, Ordering::Relaxed);
let config = ParserConfig::default().max_documents(4);
let yaml = "---\na: 1\n---\na: 2\n---\na: 3\n---\na: 4\n---\na: 5\n";
let error = parse_with_config::<CountingValue>(yaml, &config).unwrap_err();
assert!(matches!(
error,
crate::Error::Budget(crate::BudgetBreach::MaxDocuments {
limit: 4,
observed: 5
})
));
assert_eq!(DESERIALIZATIONS.load(Ordering::Relaxed), 0);
}
#[test]
fn parse_with_config_preserves_semantic_policy() {
let config = ParserConfig::default().duplicate_key_policy(crate::DuplicateKeyPolicy::Error);
let error = parse_with_config::<crate::Value>("a: 1\na: 2\n", &config).unwrap_err();
assert_eq!(error.kind(), crate::ErrorKind::DuplicateKey);
}
#[cfg_attr(
miri,
ignore = "rayon/crossbeam-epoch uses int-to-ptr casts unsupported under -Zmiri-strict-provenance"
)]
#[test]
fn parallel_errors_remain_in_source_order() {
let config = ParserConfig::default().duplicate_key_policy(crate::DuplicateKeyPolicy::Error);
let yaml = concat!(
"---\nid: 1\n",
"---\nid: 2\nid: 3\n",
"---\nid: [\n",
"---\nid: 4\n",
);
let error = parse_with_config::<crate::Value>(yaml, &config).unwrap_err();
assert_eq!(error.kind(), crate::ErrorKind::DuplicateKey);
}
#[cfg_attr(
miri,
ignore = "rayon/crossbeam-epoch uses int-to-ptr casts unsupported under -Zmiri-strict-provenance"
)]
#[test]
fn parse_propagates_first_error() {
#[derive(Debug, serde::Deserialize)]
#[allow(dead_code)]
struct Record {
id: u32,
}
let yaml = "---\nid: 1\n---\nid: [\n";
let res: Result<Vec<Record>> = parse(yaml);
assert!(res.is_err());
}
#[cfg_attr(
miri,
ignore = "rayon/crossbeam-epoch uses int-to-ptr casts unsupported under -Zmiri-strict-provenance"
)]
#[test]
fn parse_matches_sequential_for_correctness() {
let mut yaml = String::new();
for i in 0..50 {
yaml.push_str(&format!("---\nid: {i}\nname: record-{i}\n"));
}
#[derive(Debug, serde::Deserialize, PartialEq)]
struct Record {
id: u32,
name: String,
}
let parallel: Vec<Record> = parse(&yaml).unwrap();
let sequential: Vec<Record> = crate::load_all_as(&yaml).unwrap();
assert_eq!(parallel, sequential);
}
}