diff --git a/CHANGELOG.md b/CHANGELOG.md index 1a90c888d..3d02844c7 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -15,6 +15,9 @@ targets: - name: test ``` +- added two options to reverse proxy config `forced_retry_interval_secs` and `connect_timeout_secs` +`forced_retry_interval_secs` forces every x seconds a econnect to the provider, +`connect_timeout_secs` tries only x seconds for connection, if not successfull starts a retry. # 2.2.2 (2025-03-12) - !BREAKING CHANGE! Target options moved to specific target output definitions. diff --git a/README.md b/README.md index ce40524c1..efa71bc4e 100644 --- a/README.md +++ b/README.md @@ -6,10 +6,9 @@ **m3u-filter** is a versatile tool for processing playlists. Key capabilities include: -- Filtering, renaming, mapping, and sorting playlist entries and saving them in EXTM3U, XTREAM, or Kodi formats. +- Filtering, renaming, mapping, and sorting playlist entries and saving them in EXTM3U, XTREAM, or Strm (Kodi) formats. - Process multiple input files and create multiple output files through target definitions. -- Act as a simple Xtream or M3U server after processing entries. -- Serve as a redirect or reverse proxy for Xtream. +- Act as a simple Xtream or M3U redirect or reverse proxy after processing entries. - Schedule updates in server mode. - Running as a CLI tool to deliver playlists through web servers (e.g., Nginx, Apache). - Define multiple filtering targets to create several playlists from a large one. @@ -24,21 +23,18 @@ - Define HdHomeRun devices to use with Plex/Emby/Jellyfin If you need to exclude certain entries from a playlist, you can create filters using headers and apply regex-based renaming or mapping. -Run `m3u-filter` as a CLI or in server mode for a web-based UI to manage playlist content and save filtered groups. -From the Web-UI, you can view the playlist contents, filter or search entries. +Run `m3u-filter` in server mode for a web-based UI to view the playlist contents, filter or search entries. +Playlist User can also defie their Group filter through the Web-UI. ![m3u-filter_function](https://github.com/user-attachments/assets/1b5ba462-712a-4f41-9140-8cca913ba5f4) ## Starting in server mode for Web-UI The Web-UI is available in server mode. You need to start `m3u-filter` with the `-s` (`--server`) option. On the first page you can select one of the defined input sources in the configuration, or write an url to the text field. -The contents of the playlist are displayed in the tree-view. Each link has one or more buttons. -The first is for copying the url into clipboard. The others are visible if you have configured the `video`section. +The contents of the playlist are displayed in Gallery or Tree-View. Each link has one or more buttons. +The first is for copying the url into clipboard. The others are visible if you have configured the `video` section. Based on the stream type, you will be able to download or search in a configured movie database for this entry. -In the tree-view each entry has a checkbox in front. Selecting the checkbox means **discarding** this entry from the -manual download when you hit the `Save` button. - ## Command line Arguments ``` Usage: m3u-filter [OPTIONS] @@ -60,7 +56,7 @@ Options: ## 1. `config.yml` -For running in cli mode, you need to define a `config.yml` file which can be xonfig directory next to the executable or provided with the +For running in cli mode, you need to define a `config.yml` file which can be inside config directory next to the executable or provided with the `-c` cli argument. For running specific targets use the `-t` argument like `m3u-filter -t -t `. @@ -82,11 +78,13 @@ Top level entries in the config files are: * `log` _optional * `user_access_control` _optional_ * `channel_unavailable_file` _optional_ +* `hdhomerun` _optional_ ### 1.1. `threads` If you are running on a cpu which has multiple cores, you can set for example `threads: 2` to run two threads. Don't use too many threads, you should consider max of `cpu cores * 2`. Default is `0`. +If you process the same provider multiple times each thread uses a connection. Keep in mind that you hit the provider max-connection. ### 1.2. `api` `api` contains the `server-mode` settings. To run `m3u-filter` in `server-mode` you need to start it with the `-s`cli argument. @@ -98,9 +96,11 @@ Default is `0`. With this configuration, you should create a `data` directory where you execute the binary. +Be aware that different configurations (e.g. user bouquets) along the playlists are stored in this directory. + ### 1.4 `messaging` `messaging` is an optional configuration for receiving messages. -Currently only and rest is supported. +Currently `telegram`, `rest` and `pushover.net` is supported. Messaging is Opt-In, you need to set the `notify_on` message types which are - `info` diff --git a/src/api/api_utils.rs b/src/api/api_utils.rs index 94d519380..ad5cfa00c 100644 --- a/src/api/api_utils.rs +++ b/src/api/api_utils.rs @@ -120,21 +120,21 @@ pub async fn get_user_target<'a>(api_req: &'a UserApiRequest, app_state: &'a App get_user_target_by_credentials(username, password, api_req, app_state).await } -fn get_stream_options(app_state: &AppState) -> (bool, bool, usize, bool) { - let (stream_retry, buffer_enabled, buffer_size) = app_state +fn get_stream_options(app_state: &AppState) -> (bool, u32, u32, bool, usize, bool) { + let (stream_retry, stream_force_retry_secs, stream_connect_timeout_secs, buffer_enabled, buffer_size) = app_state .config .reverse_proxy .as_ref() .and_then(|reverse_proxy| reverse_proxy.stream.as_ref()) - .map_or((false, false, 0), |stream| { + .map_or((false, 0, 0, false, 0), |stream| { let (buffer_enabled, buffer_size) = stream .buffer .as_ref() .map_or((false, 0), |buffer| (buffer.enabled, buffer.size)); - (stream.retry, buffer_enabled, buffer_size) + (stream.retry, stream.forced_retry_interval_secs, stream.connect_timeout_secs, buffer_enabled, buffer_size) }); let pipe_provider_stream = !stream_retry && !buffer_enabled; - (stream_retry, buffer_enabled, buffer_size, pipe_provider_stream) + (stream_retry, stream_force_retry_secs, stream_connect_timeout_secs, buffer_enabled, buffer_size, pipe_provider_stream) } // fn get_stream_content_length(provider_response: Option<&(Vec<(String, String)>, StatusCode)>) -> u64 { @@ -160,7 +160,7 @@ pub async fn stream_response(app_state: &AppState, stream_url: &str, } } - let (stream_retry, buffer_enabled, buffer_size, direct_pipe_provider_stream) = + let (stream_retry, stream_force_retry_secs, stream_connect_timeout, buffer_enabled, buffer_size, direct_pipe_provider_stream) = get_stream_options(app_state); if let Ok(url) = Url::parse(stream_url) { @@ -168,7 +168,7 @@ pub async fn stream_response(app_state: &AppState, stream_url: &str, let (stream_opt, provider_response) = if direct_pipe_provider_stream { provider_stream::get_provider_pipe_stream(&app_state.config, &app_state.http_client, &url, req_headers, input, item_type).await } else { - let buffer_stream_options = BufferStreamOptions::new(item_type, stream_retry, buffer_enabled, buffer_size, share_stream); + let buffer_stream_options = BufferStreamOptions::new(item_type, stream_retry, stream_force_retry_secs, stream_connect_timeout, buffer_enabled, buffer_size, share_stream); provider_stream::get_provider_reconnect_buffered_stream(&app_state.config, &app_state.http_client, &url, req_headers, input, buffer_stream_options).await }; if let Some(stream) = stream_opt { diff --git a/src/api/main_api.rs b/src/api/main_api.rs index 6fbe4aecb..fbc48e4a5 100644 --- a/src/api/main_api.rs +++ b/src/api/main_api.rs @@ -51,7 +51,7 @@ async fn create_healthcheck(app_state: &Arc) -> Healthcheck { (active_user.active_users().await, active_user.active_connections().await) }; - let active_provider_connections = app_state.active_provider.active_connections().await; + let active_provider_connections = app_state.active_provider.active_connections(); let build_time: Option = BUILD_TIMESTAMP.to_string().parse::>().ok().map(|datetime| datetime.format("%Y-%m-%d %H:%M:%S %Z").to_string()); Healthcheck { diff --git a/src/api/model/active_provider_manager.rs b/src/api/model/active_provider_manager.rs index 4c225a907..55e3e500f 100644 --- a/src/api/model/active_provider_manager.rs +++ b/src/api/model/active_provider_manager.rs @@ -1,7 +1,7 @@ +use std::cell::RefCell; use crate::model::config::{Config, ConfigInput, ConfigInputAlias, InputType}; use std::collections::HashMap; -use std::sync::atomic::{AtomicUsize, Ordering}; -use tokio::sync::RwLock; +use std::sync::atomic::{AtomicU16, AtomicUsize, Ordering}; /// This struct represents an individual provider configuration with fields like: /// @@ -20,7 +20,7 @@ pub struct ProviderConfig { pub input_type: InputType, max_connections: u16, priority: i16, - current_connections: RwLock, + current_connections: AtomicU16, } impl ProviderConfig { @@ -34,7 +34,7 @@ impl ProviderConfig { input_type: cfg.input_type.clone(), max_connections: cfg.max_connections, priority: cfg.priority, - current_connections: RwLock::new(0), + current_connections: AtomicU16::new(0), } } @@ -48,13 +48,13 @@ impl ProviderConfig { input_type: cfg.input_type.clone(), max_connections: alias.max_connections, priority: alias.priority, - current_connections: RwLock::new(0), + current_connections: AtomicU16::new(0), } } #[inline] - pub async fn is_exhausted(&self) -> bool { - self.max_connections > 0 && *self.current_connections.read().await >= self.max_connections + pub fn is_exhausted(&self) -> bool { + self.max_connections > 0 && self.current_connections.load(Ordering::SeqCst) >= self.max_connections } // // #[inline] @@ -62,24 +62,24 @@ impl ProviderConfig { // !self.is_exhausted() // } - pub async fn try_allocate(&self, force: bool) -> bool { - let mut connections = self.current_connections.write().await; - if force || *connections < self.max_connections { - *connections += 1; + pub fn try_allocate(&self, force: bool) -> bool { + let connections = self.current_connections.load(Ordering::SeqCst); + if force || self.max_connections == 0 || connections < self.max_connections { + self.current_connections.fetch_add(1, Ordering::SeqCst); return true; } false } - pub async fn release(&self) { - let mut connections = self.current_connections.write().await; - if *connections > 0 { - *connections -= 1; + pub fn release(&self) { + let connections = self.current_connections.load(Ordering::SeqCst); + if connections > 0 { + self.current_connections.fetch_sub(1, Ordering::SeqCst); } } - pub async fn get_connection(&self) -> u16 { - *self.current_connections.read().await + pub fn get_connection(&self) -> u16 { + self.current_connections.load(Ordering::SeqCst) } } @@ -94,17 +94,17 @@ enum ProviderLineup { } impl ProviderLineup { - async fn acquire(&self, force: bool) -> Option<&ProviderConfig> { + fn acquire(&self, force: bool) -> Option<&ProviderConfig> { match self { - ProviderLineup::Single(lineup) => lineup.acquire(force).await, - ProviderLineup::Multi(lineup) => lineup.acquire(force).await, + ProviderLineup::Single(lineup) => lineup.acquire(force), + ProviderLineup::Multi(lineup) => lineup.acquire(force), } } - async fn release(&self, provider_id: u16) { + fn release(&self, provider_name: &str) { match self { - ProviderLineup::Single(lineup) => lineup.release(provider_id).await, - ProviderLineup::Multi(lineup) => lineup.release(provider_id).await, + ProviderLineup::Single(lineup) => lineup.release(provider_name), + ProviderLineup::Multi(lineup) => lineup.release(provider_name), } } } @@ -122,17 +122,17 @@ impl SingleProviderLineup { } } - async fn acquire(&self, force: bool) -> Option<&ProviderConfig> { - if self.provider.try_allocate(force).await { + fn acquire(&self, force: bool) -> Option<&ProviderConfig> { + if self.provider.try_allocate(force) { Some(&self.provider) } else { None } } - async fn release(&self, provider_id: u16) { - if self.provider.id == provider_id { - self.provider.release().await; + fn release(&self, provider_name: &str) { + if self.provider.name == provider_name { + self.provider.release(); } } } @@ -149,12 +149,12 @@ enum ProviderPriorityGroup { } impl ProviderPriorityGroup { - async fn is_exhausted(&self) -> bool { + fn is_exhausted(&self) -> bool { match self { - ProviderPriorityGroup::SingleProviderGroup(g) => g.is_exhausted().await, + ProviderPriorityGroup::SingleProviderGroup(g) => g.is_exhausted(), ProviderPriorityGroup::MultiProviderGroup(_, groups) => { for g in groups { - if !g.is_exhausted().await { + if !g.is_exhausted() { return false; } } @@ -232,10 +232,10 @@ impl MultiProviderLineup { /// println!("No available providers in group 0."); /// } /// ``` - async fn acquire_next_provider_from_group(priority_group: &ProviderPriorityGroup) -> Option<&ProviderConfig> { + fn acquire_next_provider_from_group(priority_group: &ProviderPriorityGroup) -> Option<&ProviderConfig> { match priority_group { ProviderPriorityGroup::SingleProviderGroup(p) => { - if p.try_allocate(false).await { + if p.try_allocate(false) { return Some(p); } } @@ -245,7 +245,7 @@ impl MultiProviderLineup { for _ in 0..provider_count { let p = pg.get(idx).unwrap(); idx = (idx + 1) % provider_count; - if p.try_allocate(false).await { + if p.try_allocate(false) { index.store(idx, Ordering::SeqCst); return Some(p); } @@ -286,15 +286,15 @@ impl MultiProviderLineup { /// println!("No available providers."); /// } /// ``` - async fn acquire(&self, force: bool) -> Option<&ProviderConfig> { + fn acquire(&self, force: bool) -> Option<&ProviderConfig> { let mut main_idx = self.index.load(Ordering::SeqCst); let provider_count = self.providers.len(); for _ in 0..provider_count { let priority_group = &self.providers[main_idx]; main_idx = (main_idx + 1) % provider_count; - if let Some(provider) = Self::acquire_next_provider_from_group(priority_group).await { - if priority_group.is_exhausted().await { + if let Some(provider) = Self::acquire_next_provider_from_group(priority_group) { + if priority_group.is_exhausted() { self.index.store(main_idx, Ordering::SeqCst); } return Some(provider); @@ -319,19 +319,19 @@ impl MultiProviderLineup { } - async fn release(&self, provider_id: u16) { + fn release(&self, provider_name: &str) { for g in &self.providers { match g { ProviderPriorityGroup::SingleProviderGroup(pc) => { - if pc.id == provider_id { - pc.release().await; + if pc.name == provider_name { + pc.release(); break; } } ProviderPriorityGroup::MultiProviderGroup(_, group) => { for pc in group { - if pc.id == provider_id { - pc.release().await; + if pc.name == provider_name { + pc.release(); return; } } @@ -344,7 +344,7 @@ impl MultiProviderLineup { pub struct ActiveProviderManager { user_access_control: bool, - providers: HashMap, + providers: Vec, } impl ActiveProviderManager { @@ -352,7 +352,7 @@ impl ActiveProviderManager { let user_access_control = cfg.user_access_control; let mut this = Self { user_access_control, - providers: HashMap::new(), + providers: Vec::new(), }; for source in &cfg.sources { for input in &source.inputs { @@ -368,41 +368,28 @@ impl ActiveProviderManager { } else { ProviderLineup::Single(SingleProviderLineup::new(input)) }; - self.providers.insert(input.name.to_string(), lineup); + self.providers.push(lineup); } - pub async fn acquire_connection(&self, lineup_name: &str) -> Option<&ProviderConfig> { - match self.providers.get(lineup_name) { - None => None, - Some(lineup) => lineup.acquire(self.user_access_control).await - } - } - - pub async fn release_connection(&self, lineup_name: &str, provider_id: u16) { - if let Some(lineup) = self.providers.get(lineup_name) { - lineup.release(provider_id).await; - } - } - - pub async fn active_connections(&self) -> Option> { - let mut result = HashMap::::new(); - for lineup in self.providers.values() { - match lineup { - ProviderLineup::Single(provider_lineup) => { - let count = *provider_lineup.provider.current_connections.read().await; - result.insert(provider_lineup.provider.name.to_string(), count); + fn get_provider_config(&self, name: &str) -> Option<(&ProviderLineup, &ProviderConfig)> { + for lineup in &self.providers { + match lineup { + ProviderLineup::Single(single) => { + if single.provider.name == name { + return Some((lineup, &single.provider)); + } } - ProviderLineup::Multi(provider_lineup) => { - for provider_group in &provider_lineup.providers { - match provider_group { - ProviderPriorityGroup::SingleProviderGroup(provider) => { - let count = *provider.current_connections.read().await; - result.insert(provider.name.to_string(), count); + ProviderLineup::Multi(multi) => { + for group in &multi.providers { + match group { + ProviderPriorityGroup::SingleProviderGroup(config) => { + return Some((lineup, config)); } - ProviderPriorityGroup::MultiProviderGroup(_, providers) => { - for provider in providers { - let count = *provider.current_connections.read().await; - result.insert(provider.name.to_string(), count); + ProviderPriorityGroup::MultiProviderGroup(_, configs) => { + for config in configs { + if config.name == name { + return Some((lineup, config)); + } } } } @@ -410,10 +397,56 @@ impl ActiveProviderManager { } } } - if result.is_empty() { + None + } + + pub fn acquire_connection(&self, input_name: &str) -> Option<&ProviderConfig> { + match self.get_provider_config(input_name) { + None => None, + Some((lineup, _config)) => lineup.acquire(self.user_access_control) + } + } + + pub fn release_connection(&self, provider_name: &str) { + if let Some((lineup, _config)) = self.get_provider_config(provider_name) { + lineup.release(provider_name); + } + } + + pub fn active_connections(&self) -> Option> { + let result = RefCell::new(HashMap::::new()); + let add_provider = |provider: &ProviderConfig| { + let count = provider.current_connections.load(Ordering::SeqCst); + if count > 0 { + result.borrow_mut().insert(provider.name.to_string(), count); + } + }; + for lineup in &self.providers { + match lineup { + ProviderLineup::Single(provider_lineup) => { + add_provider(&provider_lineup.provider); + } + ProviderLineup::Multi(provider_lineup) => { + for provider_group in &provider_lineup.providers { + match provider_group { + ProviderPriorityGroup::SingleProviderGroup(provider) => { + add_provider(provider); + } + ProviderPriorityGroup::MultiProviderGroup(_, providers) => { + for provider in providers { + add_provider(provider); + } + } + } + } + } + } + } + let status = result.take(); + if status.is_empty() { None } else { - Some(result) + Some(status) } } } diff --git a/src/api/model/event_manager.rs b/src/api/model/event_manager.rs index 4c1fd18f0..d2e7a5a04 100644 --- a/src/api/model/event_manager.rs +++ b/src/api/model/event_manager.rs @@ -38,7 +38,7 @@ impl EventManager { } if let Some(input) = input_name { // TODO this is the wrong place, move it later to the right place - // self.active_provider.acquire_connection(&input).await; + self.active_provider.acquire_connection(&input); } } Event::StreamDisconnect((username, input_name)) => { @@ -47,7 +47,7 @@ impl EventManager { info!("Active clients: {client_count}, active connections {connection_count}"); } if let Some(input) = input_name { - // self.active_provider.release_connection(&input).await; + self.active_provider.release_connection(&input); } } }; diff --git a/src/api/model/streams/mod.rs b/src/api/model/streams/mod.rs index e290ed13d..886de44e4 100644 --- a/src/api/model/streams/mod.rs +++ b/src/api/model/streams/mod.rs @@ -3,6 +3,7 @@ pub(in crate::api) mod persist_pipe_stream; pub(in crate::api) mod provider_stream_factory; pub(in crate::api) mod shared_stream_manager; pub(in crate::api) mod active_client_stream; +mod timed_client_stream; mod buffered_stream; mod client_stream; mod freeze_frame_stream; \ No newline at end of file diff --git a/src/api/model/streams/provider_stream_factory.rs b/src/api/model/streams/provider_stream_factory.rs index 17665c30d..43c2c993f 100644 --- a/src/api/model/streams/provider_stream_factory.rs +++ b/src/api/model/streams/provider_stream_factory.rs @@ -20,6 +20,7 @@ use std::sync::atomic::{AtomicUsize, Ordering}; use std::sync::Arc; use std::time::{Duration, Instant}; use url::Url; +use crate::api::model::streams::timed_client_stream::TimedClientStream; // TODO make this configurable pub const STREAM_QUEUE_SIZE: usize = 4096; // mpsc channel holding messages. with possible 8092byte chunks @@ -32,6 +33,8 @@ pub struct BufferStreamOptions { #[allow(dead_code)] item_type: PlaylistItemType, reconnect_enabled: bool, + force_reconnect_secs: u32, + connect_timeout_secs: u32, buffer_enabled: bool, buffer_size: usize, share_stream: bool, @@ -41,6 +44,8 @@ impl BufferStreamOptions { pub(crate) fn new( item_type: PlaylistItemType, reconnect_enabled: bool, + force_reconnect_secs: u32, + connect_timeout: u32, buffer_enabled: bool, buffer_size: usize, share_stream: bool @@ -48,6 +53,8 @@ impl BufferStreamOptions { Self { item_type, reconnect_enabled, + force_reconnect_secs, + connect_timeout_secs: connect_timeout, buffer_enabled, buffer_size, share_stream, @@ -87,6 +94,8 @@ struct ProviderStreamOptions { continue_flag: Arc, url: Url, reconnect: bool, + reconnect_force_secs: u32, + connect_timeout_secs: u32, headers: HeaderMap, range_bytes: Arc>, } @@ -170,7 +179,7 @@ fn get_request_range_start_bytes(req_headers: &HashMap>) -> Opti fn get_client_stream_request_params( req_headers: &HeaderMap, input: Option<&ConfigInput>, - options: &BufferStreamOptions) -> (usize, Option, bool, HeaderMap) + options: &BufferStreamOptions) -> (usize, Option, bool, u32, u32, HeaderMap) { let stream_buffer_size = if options.is_buffer_enabled() { options.get_stream_buffer_size() } else { 1 }; let filter_header = get_header_filter_for_item_type(options.item_type); @@ -185,7 +194,7 @@ fn get_client_stream_request_params( // We merge configured input headers with the headers from the request. let headers = get_request_headers(input_headers.as_ref(), Some(&req_headers)); - (stream_buffer_size, req_range_start_bytes, options.is_reconnect_enabled(), headers) + (stream_buffer_size, req_range_start_bytes, options.is_reconnect_enabled(), options.force_reconnect_secs, options.connect_timeout_secs, headers) } fn prepare_client(request_client: &Arc, url: &Url, headers: &HeaderMap, range_start_bytes_to_request: Option) -> (reqwest::RequestBuilder, bool) { @@ -248,15 +257,25 @@ async fn stream_provider(client: Arc, stream_options: ProviderS debug_if_enabled!("stream provider {}", sanitize_sensitive_info(url.as_str())); while stream_options.should_continue() { debug_if_enabled!("Reconnecting stream {}", sanitize_sensitive_info(url.as_str())); - let (client, _) = prepare_client(&client, url, headers, range_start); + let (client_builder, _) = prepare_client(&client, url, headers, range_start); + let client = if stream_options.connect_timeout_secs > 0 { + client_builder.timeout(Duration::from_secs(u64::from(stream_options.connect_timeout_secs))) + } else { + client_builder + }; match client.send().await { Ok(response) => { let status = response.status(); if status.is_success() { - return Some(response.bytes_stream().map_err(|err| { + let provider_stream = response.bytes_stream().map_err(|err| { error!("Stream error {err}"); StreamError::reqwest(&err) - }).boxed()); + }).boxed(); + return if stream_options.reconnect_force_secs > 0 { + Some(TimedClientStream::new(provider_stream, stream_options.reconnect_force_secs).boxed()) + } else { + Some(provider_stream) + } } if status.is_client_error() { return None; @@ -319,7 +338,8 @@ fn create_provider_stream_options(stream_url: &Url, req_headers: &HeaderMap, input: Option<&ConfigInput>, options: &BufferStreamOptions) -> ProviderStreamOptions { - let (buffer_size, req_range_start_bytes, reconnect, headers) = get_client_stream_request_params(req_headers, input, options); + let (buffer_size, req_range_start_bytes, reconnect, reconnect_force_secs, connect_timeout_secs, headers) + = get_client_stream_request_params(req_headers, input, options); let url = stream_url.clone(); let range_bytes = Arc::new(req_range_start_bytes.map(AtomicUsize::new)); let continue_flag = Arc::new(AtomicOnceFlag::new()); @@ -329,6 +349,8 @@ fn create_provider_stream_options(stream_url: &Url, continue_flag, url, reconnect, + reconnect_force_secs, + connect_timeout_secs, headers, range_bytes, } diff --git a/src/api/model/streams/timed_client_stream.rs b/src/api/model/streams/timed_client_stream.rs new file mode 100644 index 000000000..8d98c4339 --- /dev/null +++ b/src/api/model/streams/timed_client_stream.rs @@ -0,0 +1,29 @@ +use crate::api::model::stream_error::StreamError; +use crate::api::model::streams::provider_stream_factory::ResponseStream; +use bytes::Bytes; +use futures::Stream; +use std::pin::Pin; +use std::task::Poll; +use std::time::{Duration, Instant}; + +pub struct TimedClientStream { + inner: ResponseStream, + duration: Duration, + start_time: Instant, +} + +impl TimedClientStream { + pub(crate) fn new(inner: ResponseStream, duration: u32) -> Self { + Self { inner, duration: Duration::from_secs(u64::from(duration)) , start_time: Instant::now() } + } +} +impl Stream for TimedClientStream { + type Item = Result; + + fn poll_next(mut self: Pin<&mut Self>,cx: &mut std::task::Context<'_>,) -> Poll> { + if self.start_time.elapsed() > self.duration { + return Poll::Ready(None); + } + Pin::as_mut(&mut self.inner).poll_next(cx) + } +} \ No newline at end of file diff --git a/src/model/config.rs b/src/model/config.rs index ac7216fa4..d99717a96 100644 --- a/src/model/config.rs +++ b/src/model/config.rs @@ -1108,6 +1108,11 @@ pub struct StreamConfig { pub retry: bool, #[serde(default, skip_serializing_if = "Option::is_none")] pub buffer: Option, + #[serde(default)] + pub forced_retry_interval_secs: u32, + #[serde(default)] + pub connect_timeout_secs: u32, + } impl StreamConfig {