diff --git a/src/api/model/buffered_stream.rs b/src/api/model/buffered_stream.rs index 4149f583e..ed52f6270 100644 --- a/src/api/model/buffered_stream.rs +++ b/src/api/model/buffered_stream.rs @@ -19,7 +19,6 @@ use tokio_stream::Stream; use url::Url; const STREAM_QUEUE_SIZE: usize = 1024; // mpsc channel holding messages. -const STREAM_CONNECT_TIMEOUT_SECS: u64 = 5; // Wait timeout secs for connection when server dropped connection, then retry const ERR_RETRY_TIMEOUT_SECS: u64 = 10; // If connect status is 4xx or 5xx, we wait until we allow next request from client fn get_request_bytes(req_headers: &HashMap>) -> usize { @@ -78,13 +77,13 @@ pub fn get_buffered_stream(stream_url: &Url, req: &HttpRequest, input: Option<&C let url = stream_url.clone(); let stop_signal = Arc::new(AtomicBool::new(false)); let stop_stream = Arc::clone(&stop_signal); + let req_client = request_utils::get_client_request(input_headers.as_ref(), &url, Some(&req_headers)); actix_rt::spawn(async move { let masked_url = mask_sensitive_info(url.as_str()); let req_bytes = get_request_bytes(&req_headers); let bytes_counter = if range_send { Some(AtomicUsize::new(req_bytes)) } else { None }; while !stop_signal.load(Ordering::Relaxed) { - let mut client = request_utils::get_client_request(input_headers.as_ref(), &url, Some(&req_headers)); - client = client.timeout(Duration::from_secs(STREAM_CONNECT_TIMEOUT_SECS)); + let Some(mut client) = req_client.try_clone() else { break }; let bytes_to_request = bytes_counter.as_ref().map_or(0, |atomic| atomic.load(Ordering::Relaxed)); if bytes_to_request > 0 { // on reconnect send range header to avoid starting from beginning for vod @@ -108,7 +107,7 @@ pub fn get_buffered_stream(stream_url: &Url, req: &HttpRequest, input: Option<&C match response.chunk().await { Ok(Some(chunk)) => { if chunk.is_empty() { - // debug!("Stream finished ? {masked_url}"); + // debug!("Download Stream finished ? {masked_url}"); break; } if let Ok(permit) = tx.reserve().await { @@ -118,18 +117,19 @@ pub fn get_buffered_stream(stream_url: &Url, req: &HttpRequest, input: Option<&C bytes.fetch_add(len, Ordering::Relaxed); } } else { - // debug!("Stream finished, client disconnect ? {masked_url}"); + // debug!("Client disconnect ? {masked_url}"); stop_signal.store(true, Ordering::Relaxed); break; } } Err(_err) => { + // debug!("Media stream error {masked_url} {_err:?}"); stop_signal.store(true, Ordering::Relaxed); break; } Ok(None) => { // no chunk available - // debug!("media stream finished no data available"); + // debug!("Media stream finished no data available {masked_url}"); stop_signal.store(true, Ordering::Relaxed); break; } @@ -149,7 +149,7 @@ pub fn get_buffered_stream(stream_url: &Url, req: &HttpRequest, input: Option<&C } actix_web::rt::time::sleep(Duration::from_secs(1)).await; } - debug!("Reconnecting stream stopped {masked_url}"); + // debug!("Reconnecting stream stopped {masked_url}"); drop(tx); }); diff --git a/src/api/xmltv_api.rs b/src/api/xmltv_api.rs index 9b8ae8e2e..26d0c95b1 100644 --- a/src/api/xmltv_api.rs +++ b/src/api/xmltv_api.rs @@ -121,10 +121,10 @@ async fn serve_epg(epg_path: &Path, req: &HttpRequest, user: &ProxyUserCredentia fn serve_epg_with_timeshift(epg_file: File, offset_minutes: i32) -> HttpResponse { let reader = BufReader::new(epg_file); - let encoder = GzEncoder::new(Vec::new(), Compression::default()); + let encoder = GzEncoder::new(Vec::with_capacity(4096), Compression::default()); let mut xml_reader = Reader::from_reader(reader); let mut xml_writer = Writer::new(encoder); - let mut buf = Vec::new(); + let mut buf = Vec::with_capacity(1024); let duration = Duration::minutes(i64::from(offset_minutes)); loop { diff --git a/src/filter.rs b/src/filter.rs index 1ae3747cc..ac7938d87 100644 --- a/src/filter.rs +++ b/src/filter.rs @@ -319,7 +319,7 @@ macro_rules! handle_expr { } fn get_parser_expression(expr: Pair, templates: &Vec, errors: &mut Vec) -> Filter { - let mut stmts = Vec::new(); + let mut stmts = Vec::with_capacity(128); let pairs = expr.into_inner(); let mut bop: Option = None; let mut uop: Option = None; @@ -383,7 +383,7 @@ fn get_parser_binary_op(expr: &Pair) -> Result>) -> Result { - let empty_list = Vec::new(); + let empty_list = Vec::with_capacity(0); let template_list: &Vec = templates.unwrap_or(&empty_list); let source = apply_templates_to_pattern(filter_text, template_list); diff --git a/src/model/config.rs b/src/model/config.rs index ed32fd145..280928848 100644 --- a/src/model/config.rs +++ b/src/model/config.rs @@ -949,7 +949,7 @@ impl Config { for source in &mut self.sources { for target in &mut source.targets { if let Some(mapping_ids) = &target.mapping { - let mut target_mappings = Vec::new(); + let mut target_mappings = Vec::with_capacity(128); for mapping_id in mapping_ids { let mapping = mappings_cfg.get_mapping(mapping_id); if let Some(mappings) = mapping { diff --git a/src/model/mapping.rs b/src/model/mapping.rs index a10251559..e67a41848 100644 --- a/src/model/mapping.rs +++ b/src/model/mapping.rs @@ -308,11 +308,12 @@ impl MappingValueProcessor<'_> { .map(|caps| caps.as_str()) .collect::>(); + let mut captured_tag_values: Vec<&str> = Vec::with_capacity(128); for tag_capture in tag_captures { for mapping_tag in &self.mapper.t_tags { if mapping_tag.name.eq(tag_capture) { // we have the right tag, now get all captured values - let mut captured_tag_values: Vec<&str> = Vec::new(); + captured_tag_values.clear(); for cap in &mapping_tag.captures { if let Some(cap_value) = captures.get(cap.as_str()) { captured_tag_values.push(cap_value); diff --git a/src/processing/playlist_processor.rs b/src/processing/playlist_processor.rs index 94841962d..711090d1e 100644 --- a/src/processing/playlist_processor.rs +++ b/src/processing/playlist_processor.rs @@ -42,7 +42,7 @@ fn is_valid(pli: &PlaylistItem, target: &ConfigTarget) -> bool { #[allow(clippy::unnecessary_wraps)] fn filter_playlist(playlist: &mut [PlaylistGroup], target: &ConfigTarget) -> Option> { debug!("Filtering {} groups", playlist.len()); - let mut new_playlist = Vec::new(); + let mut new_playlist = Vec::with_capacity(128); for pg in playlist.iter_mut() { let channels = pg.channels.iter() .filter(|&pli| is_valid(pli, target)).cloned().collect::>(); @@ -151,7 +151,7 @@ fn rename_playlist(playlist: &mut [PlaylistGroup], target: &ConfigTarget) -> Opt match &target.rename { Some(renames) => { if !renames.is_empty() { - let mut new_playlist: Vec = Vec::new(); + let mut new_playlist: Vec = Vec::with_capacity(playlist.len()); for g in playlist { let mut grp = g.clone(); for r in renames { @@ -219,7 +219,7 @@ fn map_playlist(playlist: &mut [PlaylistGroup], target: &ConfigTarget) -> Option // if the group names are changed, restructure channels to the right groups // we use - let mut new_groups: Vec = Vec::new(); + let mut new_groups: Vec = Vec::with_capacity(128); let mut grp_id: u32 = 0; for playlist_group in new_playlist { for channel in &playlist_group.channels { @@ -296,7 +296,7 @@ async fn process_source(cfg: Arc, source_idx: usize, user_targets: Arc

::new(); let mut target_stats = Vec::::new(); - let mut source_playlists = Vec::new(); + let mut source_playlists = Vec::with_capacity(128); let enabled_inputs = source.inputs.iter().filter(|item| item.enabled).count(); // Downlod the sources for input in &source.inputs { diff --git a/src/repository/bplustree.rs b/src/repository/bplustree.rs index 7240fc4d5..e5228abe5 100644 --- a/src/repository/bplustree.rs +++ b/src/repository/bplustree.rs @@ -318,8 +318,7 @@ where let bytes_available_on_block = BLOCK_SIZE - read_pos; let content_bytes = if values_length > bytes_available_on_block { let mut left_over_bytes = values_length; - let mut content_chunk = Vec::new(); - content_chunk.extend(&buffer[read_pos..read_pos + bytes_available_on_block]); + let mut content_chunk = Vec::from(&buffer[read_pos..read_pos + bytes_available_on_block]); left_over_bytes -= bytes_available_on_block; while left_over_bytes > 0 { file.read_exact(buffer)?; @@ -377,14 +376,14 @@ where fn decode_content(content_bytes: &Vec) -> Option> { if let Ok(mut decoder) = StreamingDecoder::new(&**content_bytes) { - let mut result = Vec::new(); + let mut result = Vec::with_capacity(content_bytes.len()); if decoder.read_to_end(&mut result).is_ok() { return Some(result) } } // TODO remove at next deployment, this is only fallback for older compressed files - let mut decoder = flate2::write::ZlibDecoder::new(Vec::new()); + let mut decoder = flate2::write::ZlibDecoder::new(Vec::with_capacity(content_bytes.len())); if let Ok(()) = decoder.write_all(content_bytes) { if let Ok(decoded) = decoder.finish() { return Some(decoded); diff --git a/src/repository/indexed_document.rs b/src/repository/indexed_document.rs index c26289a18..b9a1b3800 100644 --- a/src/repository/indexed_document.rs +++ b/src/repository/indexed_document.rs @@ -170,7 +170,8 @@ where if size == encoded_bytes.len() { // check if it is equal let mut record_buffer = vec![0; size]; - record_buffer.resize(size, 0); + // record_buffer.resize(size, 0); + self.main_file.read_exact(&mut record_buffer)?; if record_buffer == encoded_bytes { return Ok(()); @@ -310,7 +311,7 @@ where offsets, index: 0, failed: false, - t_buffer: Vec::new(), + t_buffer: Vec::with_capacity(BLOCK_SIZE), t_type: PhantomData, k_type: PhantomData, }) diff --git a/src/utils/directed_graph.rs b/src/utils/directed_graph.rs index 0e982f9e2..17de9006b 100644 --- a/src/utils/directed_graph.rs +++ b/src/utils/directed_graph.rs @@ -42,8 +42,8 @@ where // Detect and return cycles in the graph pub fn find_cycles(&self) -> Vec> { let mut visited = HashSet::new(); - let mut recursion_stack = Vec::new(); - let mut cycles = Vec::new(); + let mut recursion_stack = Vec::with_capacity(128); + let mut cycles = Vec::with_capacity(128); for node in self.adjacencies.keys() { if !visited.contains(node) { diff --git a/src/utils/download.rs b/src/utils/download.rs index e0903da80..b3ddc1a28 100644 --- a/src/utils/download.rs +++ b/src/utils/download.rs @@ -142,7 +142,7 @@ const ACTIONS: [(XtreamCluster, &str, &str); 3] = [ (XtreamCluster::Series, "get_series_categories", "get_series")]; pub async fn get_xtream_playlist(input: &ConfigInput, working_dir: &str) -> (Vec, Vec) { - let mut playlist_groups: Vec = Vec::new(); + let mut playlist_groups: Vec = Vec::with_capacity(128); let username = input.username.as_ref().map_or("", |v| v); let password = input.password.as_ref().map_or("", |v| v); let base_url = format!("{}/player_api.php?username={}&password={}", input.url, username, password); diff --git a/src/utils/json_utils.rs b/src/utils/json_utils.rs index 7e2ad6692..cfdbcaa38 100644 --- a/src/utils/json_utils.rs +++ b/src/utils/json_utils.rs @@ -64,7 +64,7 @@ pub fn json_iter_array( } pub fn json_filter_file(file_path: &Path, filter: &HashMap<&str, &str>) -> Vec { - let mut filtered: Vec = Vec::new(); + let mut filtered: Vec = Vec::with_capacity(1024); if !file_path.exists() { return filtered; // Return early if the file does not exist } diff --git a/src/utils/multi_file_reader.rs b/src/utils/multi_file_reader.rs index dcbba2cf6..d4f013769 100644 --- a/src/utils/multi_file_reader.rs +++ b/src/utils/multi_file_reader.rs @@ -9,7 +9,7 @@ pub struct MultiFileReader { impl MultiFileReader { pub fn new(paths: &Vec) -> io::Result { - let mut files = Vec::new(); + let mut files = Vec::with_capacity(paths.len()); for path in paths { match File::open(path) { Ok(file) => { files.push(file); }