From 98dc80dd758ff4439b03e6ec7fa0da44e8b63e8a Mon Sep 17 00:00:00 2001 From: euzu Date: Fri, 21 Nov 2025 01:35:04 +0100 Subject: [PATCH] SyncIOBridge caused tokio panic --- Cargo.lock | 6 +- backend/Cargo.toml | 4 +- .../src/api/endpoints/api_playlist_utils.rs | 14 --- backend/src/api/model/app_state.rs | 3 +- backend/src/model/stats.rs | 2 +- backend/src/model/xmltv.rs | 85 ++++++++++++------- backend/src/model/xtream.rs | 2 +- backend/src/repository/epg_repository.rs | 34 ++++---- backend/src/repository/playlist_repository.rs | 2 + backend/src/repository/user_repository.rs | 11 ++- backend/src/repository/xtream_repository.rs | 75 +++++++++------- backend/src/utils/json_utils.rs | 58 +++++++------ backend/src/utils/network/request.rs | 2 +- frontend/Cargo.toml | 4 +- shared/Cargo.toml | 2 +- shared/src/utils/json_utils.rs | 61 +------------ 16 files changed, 173 insertions(+), 192 deletions(-) diff --git a/Cargo.lock b/Cargo.lock index cdfaee97b..13cc292a9 100644 --- a/Cargo.lock +++ b/Cargo.lock @@ -1096,7 +1096,7 @@ dependencies = [ [[package]] name = "frontend" -version = "3.2.6" +version = "3.2.7" dependencies = [ "anyhow", "base64", @@ -3765,7 +3765,7 @@ dependencies = [ [[package]] name = "shared" -version = "3.2.6" +version = "3.2.7" dependencies = [ "base64", "bitflags 2.10.0", @@ -4314,7 +4314,7 @@ checksum = "e421abadd41a4225275504ea4d6566923418b7f05506fbc9c0fe86ba7396114b" [[package]] name = "tuliprox" -version = "3.2.6" +version = "3.2.7" dependencies = [ "arc-swap", "async-compression", diff --git a/backend/Cargo.toml b/backend/Cargo.toml index 325704a41..9a8fee034 100644 --- a/backend/Cargo.toml +++ b/backend/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "tuliprox" -version = "3.2.6" +version = "3.2.7" edition = "2021" # See more keys and their definitions at https://doc.rust-lang.org/cargo/reference/manifest.html @@ -47,7 +47,7 @@ tokio = { version = "1.48", features = ["rt-multi-thread", "parking_lot", "fs"] #console-subscriber = "0" #tracing = "0.1" #tracing-subscriber = { version = "0.3", features = ["fmt", "env-filter"] } -tokio-util = { version = "0.7", features = ["io-util"] } +tokio-util = { version = "0.7"} tempfile = "3.23" ruzstd = "0.8" filetime = "0.2" diff --git a/backend/src/api/endpoints/api_playlist_utils.rs b/backend/src/api/endpoints/api_playlist_utils.rs index bb84329da..9dfc8c5d1 100644 --- a/backend/src/api/endpoints/api_playlist_utils.rs +++ b/backend/src/api/endpoints/api_playlist_utils.rs @@ -99,20 +99,6 @@ fn group_playlist_groups_by_cluster(playlist: Vec) -> (Vec, Option), std::io::Error>) -> Option { -// if let Ok((Some(file_path), _content)) = action { -// if let Ok(content) = tokio::fs::read_to_string(&file_path).await { -// // TODO deserialize like sax parser -// if let Ok(categories) = serde_json::from_str::>(&content) { -// return serde_json::to_string(&categories).ok(); -// } -// } -// } -// None -// } - - async fn grouped_channels( cfg: &AppConfig, target: &ConfigTarget, diff --git a/backend/src/api/model/app_state.rs b/backend/src/api/model/app_state.rs index d4bb28321..146ed1c8f 100644 --- a/backend/src/api/model/app_state.rs +++ b/backend/src/api/model/app_state.rs @@ -10,7 +10,7 @@ use crate::repository::playlist_repository::load_target_into_memory_cache; use crate::tools::lru_cache::LRUResourceCache; use crate::utils::request::create_client; use arc_swap::{ArcSwap, ArcSwapOption}; -use log::error; +use log::{error, info}; use reqwest::Client; use shared::error::TuliproxError; use shared::model::UserConnectionPermission; @@ -207,6 +207,7 @@ pub fn create_cache(config: &Config) -> Option>> { }); let cache_enabled = lru_cache.is_some(); if cache_enabled { + info!("Scanning cache"); if let Some(res_cache) = lru_cache { let cache = Arc::new(Mutex::new(res_cache)); let cache_scanner = Arc::clone(&cache); diff --git a/backend/src/model/stats.rs b/backend/src/model/stats.rs index bc21b0321..fc2fb4ff3 100644 --- a/backend/src/model/stats.rs +++ b/backend/src/model/stats.rs @@ -8,7 +8,7 @@ pub fn format_elapsed_time(seconds: u64) -> String { } else { let minutes = seconds / 60; let seconds = seconds % 60; - format!("{minutes}:{seconds} mins") + format!("{minutes}:{seconds:02} mins") } } diff --git a/backend/src/model/xmltv.rs b/backend/src/model/xmltv.rs index 5b47f2779..62643ce6f 100644 --- a/backend/src/model/xmltv.rs +++ b/backend/src/model/xmltv.rs @@ -1,7 +1,6 @@ use crate::model::xmltv::XmlTagIcon::Undefined; use chrono::{Datelike, TimeZone, Utc}; use quick_xml::events::{BytesEnd, BytesStart, BytesText, Event}; -use quick_xml::{Error, Writer}; use shared::error::{TuliproxError, TuliproxErrorKind}; use shared::model::{parse_xmltv_time, EpgChannel, EpgProgramme, EpgTv, InputFetchMethod}; use std::cmp::{max, min}; @@ -9,7 +8,7 @@ use std::collections::HashMap; use std::path::{Path, PathBuf}; use std::sync::Arc; use futures::TryFutureExt; -use tokio::io::AsyncRead; +use tokio::io::{AsyncRead, AsyncWrite}; use url::Url; use shared::utils::sanitize_sensitive_info; use crate::api::model::AppState; @@ -60,28 +59,6 @@ impl XmlTag { self.attributes.as_ref().and_then(|attr| attr.get(attr_name)) } - fn write_to(&self, writer: &mut Writer) -> Result<(), Error> { - let mut elem = BytesStart::new(self.name.as_str()); - - // empty icon not processed - if self.icon == Undefined && self.name.eq(EPG_TAG_ICON) { - return Ok(()); - } - - if let Some(attribs) = self.attributes.as_ref() { - for (k, v) in attribs { elem.push_attribute((k.as_str(), v.as_str())); } - } - writer.write_event(Event::Start(elem))?; - if let Some(text) = self.value.as_ref() { - writer.write_event(Event::Text(BytesText::new(text.as_str())))?; - } - if let Some(children) = &self.children { - for child in children { - child.write_to(writer)?; - } - } - Ok(writer.write_event(Event::End(BytesEnd::new(self.name.as_str())))?) - } } @@ -94,16 +71,62 @@ pub struct Epg { } impl Epg { - pub fn write_to(&self, writer: &mut Writer) -> Result<(), quick_xml::Error> { + pub async fn write_to_async( + &self, + writer: &mut quick_xml::writer::Writer, + ) -> Result<(), quick_xml::Error> { + // Start tv-element let mut elem = BytesStart::new("tv"); - if let Some(attribs) = self.attributes.as_ref() { - for (k, v) in attribs { elem.push_attribute((k.as_str(), v.as_str())); } + if let Some(attrs) = &self.attributes { + for (k, v) in attrs { + elem.push_attribute((k.as_str(), v.as_str())); + } } - writer.write_event(Event::Start(elem))?; - for child in &self.children { - child.write_to(writer)?; + writer.write_event_async(Event::Start(elem)).await?; + + // Stack for iterative writing + // bool = End-Event written? + let mut stack: Vec<(&XmlTag, bool)> = self + .children + .iter() + .rev() + .map(|c| (c, false)) + .collect(); + + while let Some((tag, ended)) = stack.pop() { + if ended { + // End-Event + writer + .write_event_async(Event::End(BytesEnd::new(tag.name.as_str()))) + .await?; + } else { + // Start-Event for the tag + let mut elem = BytesStart::new(tag.name.as_str()); + if let Some(attrs) = &tag.attributes { + for (k, v) in attrs { + elem.push_attribute((k.as_str(), v.as_str())); + } + } + writer.write_event_async(Event::Start(elem)).await?; + + // write text + if let Some(text) = &tag.value { + writer.write_event_async(Event::Text(BytesText::new(text.as_str()))).await?; + } + + // End-Marker push + children push + stack.push((tag, true)); + if let Some(children) = &tag.children { + for child in children.iter().rev() { + stack.push((child, false)); + } + } + } } - Ok(writer.write_event(Event::End(BytesEnd::new("tv")))?) + + // write tv-end + writer.write_event_async(Event::End(BytesEnd::new("tv"))).await?; + Ok(()) } } diff --git a/backend/src/model/xtream.rs b/backend/src/model/xtream.rs index 61191fa06..550ed2638 100644 --- a/backend/src/model/xtream.rs +++ b/backend/src/model/xtream.rs @@ -617,7 +617,7 @@ pub fn rewrite_doc_urls(resource_url: Option<&String>, document: &mut Map Result<(), TuliproxError> { - let file = tokio::fs::File::create(path).await.map_err(|e| notify_err!(format!("failed to create epg file: {}", e)))?; - // problem quickxml is not async - let sync_writer = tokio_util::io::SyncIoBridge::new(file); - let mut writer = quick_xml::Writer::new(std::io::BufWriter::new(sync_writer)); +pub async fn epg_write_file(target: &ConfigTarget, epg: &Epg, path: &Path) -> Result<(), TuliproxError> { + let file = tokio::fs::File::create(path).await + .map_err(|e| notify_err!(format!("failed to create epg file: {}", e)))?; + let buf_writer = tokio::io::BufWriter::new(file); + let mut writer = quick_xml::writer::Writer::new(buf_writer); - writer.write_event(quick_xml::events::Event::Decl( - quick_xml::events::BytesDecl::new("1.0", Some("utf-8"), None) - )).map_err(|e| notify_err!(format!("failed to write XML header: {}", e)))?; + // XML Header + writer.write_event_async(quick_xml::events::Event::Decl(quick_xml::events::BytesDecl::new("1.0", Some("utf-8"), None))) + .await.map_err(|e| notify_err!(format!("failed to write XML header: {}", e)))?; - writer.write_event(quick_xml::events::Event::DocType(quick_xml::events::BytesText::new("tv SYSTEM \"xmltv.dtd\""))) - .map_err(|e| notify_err!(format!("failed to write doctype: {}", e)))?; + // DOCTYPE + writer.write_event_async(quick_xml::events::Event::DocType(quick_xml::events::BytesText::new("tv SYSTEM \"xmltv.dtd\""))) + .await.map_err(|e| notify_err!(format!("failed to write doctype: {}", e)))?; + // EPG Content + epg.write_to_async(&mut writer).await.map_err(|e| notify_err!(format!("failed to write epg: {}", e)))?; - epg.write_to(&mut writer).map_err(|e| notify_err!(format!("failed to write epg: {}", e)))?; - - writer.into_inner().flush().map_err(|e| notify_err!(format!("failed to flush epg: {}", e)))?; + let inner = writer.get_mut(); // Zugriff auf den BufWriter + inner.flush().await.map_err(|e| notify_err!(format!("failed to flush epg: {}", e)))?; debug_if_enabled!("Epg for target {} written to {}", target.name, path.to_str().unwrap_or("?")); Ok(()) diff --git a/backend/src/repository/playlist_repository.rs b/backend/src/repository/playlist_repository.rs index 5708113ae..519898ccb 100644 --- a/backend/src/repository/playlist_repository.rs +++ b/backend/src/repository/playlist_repository.rs @@ -17,6 +17,7 @@ use shared::utils::{is_dash_url, is_hls_url}; use shared::create_tuliprox_error; use std::path::Path; use std::sync::Arc; +use log::info; use crate::processing::processor::playlist::apply_filter_to_playlist; pub async fn persist_playlist(app_config: &AppConfig, playlist: &mut [PlaylistGroup], epg: Option<&Epg>, @@ -204,6 +205,7 @@ pub async fn load_playlists_into_memory_cache(app_state: &AppState) -> Result<() pub async fn load_target_into_memory_cache(app_state: &AppState, target: &Arc) { if target.use_memory_cache { + info!("Loading target {} into memory cache", target.name); for output in &target.output { match output { TargetOutput::Xtream(_) => { diff --git a/backend/src/repository/user_repository.rs b/backend/src/repository/user_repository.rs index 9e28fef7d..77f907346 100644 --- a/backend/src/repository/user_repository.rs +++ b/backend/src/repository/user_repository.rs @@ -259,8 +259,10 @@ async fn save_xtream_user_bouquet_for_target(config: &Config, target_name: &str, if let Some(bouquet_categories) = bouquet { if let Some(xtream_categories) = xtream_get_playlist_categories(config, target_name, cluster).await { - let filtered: Vec<&PlaylistXtreamCategory> = xtream_categories.iter().filter(|p| bouquet_categories.contains(&p.name)).collect(); - return json_write_documents_to_file(&bouquet_path, &filtered).await; + let filtered: Vec = xtream_categories.iter().filter(|p| bouquet_categories.contains(&p.name)).cloned().collect(); + return task::spawn_blocking(move || { + json_write_documents_to_file(&bouquet_path, &filtered) + }).await?; } } @@ -278,7 +280,10 @@ async fn save_m3u_user_bouquet_for_target(storage_path: &Path, target: TargetTyp }; match bouquet { Some(bouquet_categories) => { - json_write_documents_to_file(&bouquet_path, bouquet_categories).await?; + let categories = bouquet_categories.clone(); + task::spawn_blocking(move || { + json_write_documents_to_file(&bouquet_path, &categories) + }).await??; } None => if bouquet_path.exists() { tokio::fs::remove_file(bouquet_path).await?; diff --git a/backend/src/repository/xtream_repository.rs b/backend/src/repository/xtream_repository.rs index a8b4374b2..2b0ea3175 100644 --- a/backend/src/repository/xtream_repository.rs +++ b/backend/src/repository/xtream_repository.rs @@ -21,13 +21,14 @@ use serde::Serialize; use serde_json::{json, Map, Value}; use shared::error::{create_tuliprox_error, create_tuliprox_error_result, info_err, notify_err, str_to_io_error, to_io_error, TuliproxError, TuliproxErrorKind}; use shared::model::{PlaylistEntry, PlaylistGroup, PlaylistItem, PlaylistItemType, XtreamCluster, XtreamPlaylistItem}; -use shared::utils::{generate_playlist_uuid, get_u32_from_serde_value, hex_encode, json_iter_array}; +use shared::utils::{generate_playlist_uuid, get_u32_from_serde_value, hex_encode}; use std::collections::HashMap; use std::fs; use std::fs::File; -use std::io::{BufReader, BufWriter, Error, ErrorKind, Read, Write}; +use std::io::{BufWriter, Error, ErrorKind, Read, Write}; use std::path::{Path, PathBuf}; use std::sync::Arc; +use tokio::task; macro_rules! cant_write_result { ($path:expr, $err:expr) => { @@ -137,28 +138,38 @@ fn get_map_item_as_str(map: &serde_json::Map, key: &str) -> Optio None } -fn load_old_category_ids(path: &Path) -> (u32, HashMap) { - let mut result: HashMap = HashMap::new(); - let mut max_id: u32 = 0; - for (cluster, cat) in [(XtreamCluster::Live, storage_const::COL_CAT_LIVE), (XtreamCluster::Video, storage_const::COL_CAT_VOD), (XtreamCluster::Series, storage_const::COL_CAT_SERIES)] { - let col_path = get_collection_path(path, cat); - if col_path.exists() { - if let Ok(file) = File::open(col_path) { - let reader = file_reader(file); - for entry in json_iter_array::>(reader).flatten() { - if let Some(category_id) = entry.get(crate::model::XC_TAG_CATEGORY_ID).and_then(get_u32_from_serde_value) { - if let Value::Object(item) = entry { - if let Some(category_name) = get_map_item_as_str(&item, crate::model::XC_TAG_CATEGORY_NAME) { - result.insert(format!("{cluster}{category_name}"), category_id); - max_id = max_id.max(category_id); +async fn load_old_category_ids(path: &Path) -> (u32, HashMap) { + let old_path = path.to_path_buf(); + tokio::task::spawn_blocking(move || { + let mut result: HashMap = HashMap::new(); + let mut max_id: u32 = 0; + for (cluster, cat) in [(XtreamCluster::Live, storage_const::COL_CAT_LIVE), (XtreamCluster::Video, storage_const::COL_CAT_VOD), (XtreamCluster::Series, storage_const::COL_CAT_SERIES)] { + let col_path = get_collection_path(&old_path, cat); + if col_path.exists() { + if let Ok(file) = File::open(col_path) { + let reader = file_reader(file); + match serde_json::from_reader(reader) { + Ok(value) => { + if let Value::Array(list) = value { + for entry in list { + if let Some(category_id) = entry.get(crate::model::XC_TAG_CATEGORY_ID).and_then(get_u32_from_serde_value) { + if let Value::Object(item) = entry { + if let Some(category_name) = get_map_item_as_str(&item, crate::model::XC_TAG_CATEGORY_NAME) { + result.insert(format!("{cluster}{category_name}"), category_id); + max_id = max_id.max(category_id); + } + } + } + } } } + Err(_err) => {} } } } } - } - (max_id, result) + (max_id, result) + }).await.unwrap_or_else(|_| (0, HashMap::new())) } pub fn xtream_get_storage_path(cfg: &Config, target_name: &str) -> Option { @@ -218,7 +229,7 @@ pub async fn xtream_write_playlist( let mut vod_col = Vec::with_capacity(10_000); // preserve category_ids - let (max_cat_id, existing_cat_ids) = load_old_category_ids(&path); + let (max_cat_id, existing_cat_ids) = load_old_category_ids(&path).await; let mut cat_id_counter = max_cat_id; for plg in playlist.iter_mut() { if !&plg.channels.is_empty() { @@ -252,18 +263,24 @@ pub async fn xtream_write_playlist( } } - for (col_path, data) in [ - (get_collection_path(&path, storage_const::COL_CAT_LIVE), &cat_live_col), - (get_collection_path(&path, storage_const::COL_CAT_VOD), &cat_vod_col), - (get_collection_path(&path, storage_const::COL_CAT_SERIES), &cat_series_col), - ] { - match json_write_documents_to_file(&col_path, data).await { - Ok(()) => {} - Err(err) => { - errors.push(format!("Persisting collection failed: {}: {err}", col_path.display())); + let root_path = path.clone(); + let write_errors = task::spawn_blocking(move || { + let mut write_errors = vec![]; + for (col_path, data) in [ + (get_collection_path(&root_path, storage_const::COL_CAT_LIVE), &cat_live_col), + (get_collection_path(&root_path, storage_const::COL_CAT_VOD), &cat_vod_col), + (get_collection_path(&root_path, storage_const::COL_CAT_SERIES), &cat_series_col), + ] { + match json_write_documents_to_file(&col_path, data) { + Ok(()) => {} + Err(err) => { + write_errors.push(format!("Persisting collection failed: {}: {err}", col_path.display())); + } } } - } + write_errors + }).await.map_err(|e| notify_err!(format!("Task panicked: {}", e)))?; + errors.extend(write_errors); match write_playlists_to_file( cfg, diff --git a/backend/src/utils/json_utils.rs b/backend/src/utils/json_utils.rs index aa77e5476..4700fa8aa 100644 --- a/backend/src/utils/json_utils.rs +++ b/backend/src/utils/json_utils.rs @@ -1,12 +1,10 @@ -use std::collections::{HashMap, HashSet}; -use std::fs::File; -use std::io::{BufReader, Write}; -use std::path::Path; +use crate::utils::file_reader; use serde::Serialize; use serde_json::Value; -use shared::utils::json_iter_array; -use crate::utils::file_reader; -use tokio_util::io::SyncIoBridge; +use std::collections::{HashMap, HashSet}; +use std::fs::File; +use std::io::Write; +use std::path::Path; pub fn json_filter_file(file_path: &Path, filter: &HashMap<&str, HashSet, S>) -> Vec { let mut filtered: Vec = Vec::with_capacity(1024); @@ -19,33 +17,39 @@ pub fn json_filter_file(file_path: &Path, filter: & }; let reader = file_reader(file); - for entry in json_iter_array::>(reader).flatten() { - if let Some(item) = entry.as_object() { - if filter.iter().all(|(&key, filter_set)| { - item.get(key).is_some_and(|field_value| match field_value { - Value::String(s) => filter_set.contains(s.as_str()), - Value::Number(n) => filter_set.contains(n.as_str()), - _ => false, - }) - }) { - filtered.push(entry); + match serde_json::from_reader(reader) { + Ok(value) => { + if let Value::Array(list) = value { + for entry in list { + if let Some(item) = entry.as_object() { + if filter.iter().all(|(&key, filter_set)| { + item.get(key).is_some_and(|field_value| match field_value { + Value::String(s) => filter_set.contains(s.as_str()), + Value::Number(n) => filter_set.contains(n.as_str()), + _ => false, + }) + }) { + filtered.push(entry); + } + } + } } } + Err(_err) => {} } filtered } -pub async fn json_write_documents_to_file(file: &std::path::Path, value: &T) -> std::io::Result<()> +pub fn json_write_documents_to_file( + path: &std::path::Path, + value: &T, +) -> std::io::Result<()> where - T: ?Sized + Serialize, + T: Serialize, { - let file = tokio::fs::File::create(file).await?; - - let sync_writer = SyncIoBridge::new(file); - let mut buf_writer = std::io::BufWriter::new(sync_writer); - + let file = std::fs::File::create(path)?; + let mut buf_writer = std::io::BufWriter::new(file); serde_json::to_writer(&mut buf_writer, value)?; - buf_writer.flush()?; - Ok(()) -} \ No newline at end of file + buf_writer.flush() +} diff --git a/backend/src/utils/network/request.rs b/backend/src/utils/network/request.rs index 117f91e9d..23445230e 100644 --- a/backend/src/utils/network/request.rs +++ b/backend/src/utils/network/request.rs @@ -331,7 +331,7 @@ async fn get_remote_content(client: Arc, input: &InputSource, h let (mut stream, response_url) = get_remote_content_as_stream(client.clone(), url, input.method, Some(&headers)).await.map_err(|e| str_to_io_error(&format!("Failed to read content: {e}")))?; let mut content = String::new(); stream.read_to_string(&mut content).await.map_err(|e| str_to_io_error(&format!("Failed to read content: {e}")))?; - debug_if_enabled!("Request took:{} {}", format_elapsed_time(start_time.elapsed().as_secs()), sanitize_sensitive_info(url.as_str())); + debug_if_enabled!("Request took: {} {}", format_elapsed_time(start_time.elapsed().as_secs()), sanitize_sensitive_info(url.as_str())); Ok((content, response_url)) } diff --git a/frontend/Cargo.toml b/frontend/Cargo.toml index 80e4e42c4..aa9e39291 100644 --- a/frontend/Cargo.toml +++ b/frontend/Cargo.toml @@ -1,10 +1,10 @@ [package] name = "frontend" -version = "3.2.6" +version = "3.2.7" edition = "2021" [dependencies] -shared = { version = "3.2.6", path = "../shared" } +shared = { version = "3.2.7", path = "../shared" } chrono = "0" yew = "0.21" yew-router = "0.18" diff --git a/shared/Cargo.toml b/shared/Cargo.toml index a2df61e2a..fb516e9ed 100644 --- a/shared/Cargo.toml +++ b/shared/Cargo.toml @@ -1,6 +1,6 @@ [package] name = "shared" -version = "3.2.6" +version = "3.2.7" edition = "2021" [dependencies] diff --git a/shared/src/utils/json_utils.rs b/shared/src/utils/json_utils.rs index c8f3b387e..4e0ad8d04 100644 --- a/shared/src/utils/json_utils.rs +++ b/shared/src/utils/json_utils.rs @@ -1,66 +1,7 @@ -use std::io::{self, Read}; - -use serde::de::DeserializeOwned; use serde::{Deserialize}; -use serde_json::{self, Deserializer, Value}; +use serde_json::{self, Value}; use crate::utils::{humanize_snake_case}; -fn read_skipping_ws(mut reader: impl Read) -> io::Result { - loop { - let mut byte = 0u8; - reader.read_exact(std::slice::from_mut(&mut byte))?; - if !byte.is_ascii_whitespace() { - return Ok(byte); - } - } -} - -fn invalid_data(msg: &str) -> io::Error { - io::Error::new(io::ErrorKind::InvalidData, msg) -} - -fn deserialize_single(reader: R) -> io::Result { - let next_obj = Deserializer::from_reader(reader).into_iter::().next(); - next_obj.map_or_else( - || Err(invalid_data("premature EOF")), - |result| result.map_err(Into::into), - ) -} - -fn yield_next_obj( - mut reader: R, - at_start: &mut bool, -) -> io::Result> { - if *at_start { - match read_skipping_ws(&mut reader)? { - b',' => deserialize_single(reader).map(Some), - b']' => Ok(None), - _ => Err(invalid_data("`,` or `]` not found")), - } - } else { - *at_start = true; - if read_skipping_ws(&mut reader)? == b'[' { - // read the next char to see if the array is empty - let peek = read_skipping_ws(&mut reader)?; - if peek == b']' { - Ok(None) - } else { - deserialize_single(io::Cursor::new([peek]).chain(reader)).map(Some) - } - } else { - Err(invalid_data("`[` not found")) - } - } -} - -// https://stackoverflow.com/questions/68641157/how-can-i-stream-elements-from-inside-a-json-array-using-serde-json -pub fn json_iter_array( - mut reader: R, -) -> impl Iterator> { - let mut at_start = false; - std::iter::from_fn(move || yield_next_obj(&mut reader, &mut at_start).transpose()) -} - pub fn string_or_number_u32<'de, D>(deserializer: D) -> Result where D: serde::Deserializer<'de>,