SyncIOBridge caused tokio panic

This commit is contained in:
euzu
2025-11-21 01:35:04 +01:00
parent 24880101df
commit 98dc80dd75
16 changed files with 173 additions and 192 deletions
Generated
+3 -3
View File
@@ -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",
+2 -2
View File
@@ -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"
@@ -99,20 +99,6 @@ fn group_playlist_groups_by_cluster(playlist: Vec<PlaylistGroup>) -> (Vec<Playli
(live, video, series)
}
// async fn get_categories_content(action: Result<(Option<PathBuf>, Option<String>), std::io::Error>) -> Option<String> {
// 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::<Vec<PlaylistXtreamCategory>>(&content) {
// return serde_json::to_string(&categories).ok();
// }
// }
// }
// None
// }
async fn grouped_channels(
cfg: &AppConfig,
target: &ConfigTarget,
+2 -1
View File
@@ -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<Arc<Mutex<LRUResourceCache>>> {
});
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);
+1 -1
View File
@@ -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")
}
}
+54 -31
View File
@@ -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<W: std::io::Write>(&self, writer: &mut Writer<W>) -> 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<W: std::io::Write>(&self, writer: &mut Writer<W>) -> Result<(), quick_xml::Error> {
pub async fn write_to_async<W: AsyncWrite + Unpin>(
&self,
writer: &mut quick_xml::writer::Writer<W>,
) -> 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(())
}
}
+1 -1
View File
@@ -617,7 +617,7 @@ pub fn rewrite_doc_urls(resource_url: Option<&String>, document: &mut Map<String
}
#[derive(Deserialize, Serialize)]
#[derive(Deserialize, Serialize, Clone)]
pub struct PlaylistXtreamCategory {
#[serde(alias = "category_id")]
pub id: String,
+18 -16
View File
@@ -1,29 +1,31 @@
use std::io::Write;
use shared::error::{notify_err, TuliproxError};
use crate::model::{Config, ConfigTarget, TargetOutput};
use crate::model::Epg;
use crate::model::{Config, ConfigTarget, TargetOutput};
use crate::repository::m3u_repository::m3u_get_epg_file_path;
use crate::repository::xtream_repository::{xtream_get_epg_file_path, xtream_get_storage_path};
use crate::utils::debug_if_enabled;
use shared::error::{notify_err, TuliproxError};
use std::path::Path;
use tokio::io::AsyncWriteExt;
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)))?;
// 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<tokio::fs::File>
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(())
@@ -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<ConfigTarget>) {
if target.use_memory_cache {
info!("Loading target {} into memory cache", target.name);
for output in &target.output {
match output {
TargetOutput::Xtream(_) => {
+8 -3
View File
@@ -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<PlaylistXtreamCategory> = 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?;
+46 -29
View File
@@ -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<String, Value>, key: &str) -> Optio
None
}
fn load_old_category_ids(path: &Path) -> (u32, HashMap<String, u32>) {
let mut result: HashMap<String, u32> = 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::<Value, BufReader<File>>(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<String, u32>) {
let old_path = path.to_path_buf();
tokio::task::spawn_blocking(move || {
let mut result: HashMap<String, u32> = 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<PathBuf> {
@@ -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,
+31 -27
View File
@@ -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<S: ::std::hash::BuildHasher>(file_path: &Path, filter: &HashMap<&str, HashSet<String, S>, S>) -> Vec<serde_json::Value> {
let mut filtered: Vec<serde_json::Value> = Vec::with_capacity(1024);
@@ -19,33 +17,39 @@ pub fn json_filter_file<S: ::std::hash::BuildHasher>(file_path: &Path, filter: &
};
let reader = file_reader(file);
for entry in json_iter_array::<serde_json::Value, BufReader<File>>(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<T>(file: &std::path::Path, value: &T) -> std::io::Result<()>
pub fn json_write_documents_to_file<T>(
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(())
}
buf_writer.flush()
}
+1 -1
View File
@@ -331,7 +331,7 @@ async fn get_remote_content(client: Arc<reqwest::Client>, 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))
}
+2 -2
View File
@@ -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"
+1 -1
View File
@@ -1,6 +1,6 @@
[package]
name = "shared"
version = "3.2.6"
version = "3.2.7"
edition = "2021"
[dependencies]
+1 -60
View File
@@ -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<u8> {
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<T: DeserializeOwned, R: Read>(reader: R) -> io::Result<T> {
let next_obj = Deserializer::from_reader(reader).into_iter::<T>().next();
next_obj.map_or_else(
|| Err(invalid_data("premature EOF")),
|result| result.map_err(Into::into),
)
}
fn yield_next_obj<T: DeserializeOwned, R: Read>(
mut reader: R,
at_start: &mut bool,
) -> io::Result<Option<T>> {
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<T: DeserializeOwned, R: Read>(
mut reader: R,
) -> impl Iterator<Item = Result<T, io::Error>> {
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<u32, D::Error>
where
D: serde::Deserializer<'de>,