ActiveProvider status tracking added.

connect timeout for provider streams added
This commit is contained in:
euzu
2025-03-13 12:02:39 +01:00
parent fad4e0f421
commit 46aa7b04e5
10 changed files with 198 additions and 105 deletions
+3
View File
@@ -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.
+12 -12
View File
@@ -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 <target_name> -t <other_target_name>`.
@@ -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`
+7 -7
View File
@@ -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 {
+1 -1
View File
@@ -51,7 +51,7 @@ async fn create_healthcheck(app_state: &Arc<AppState>) -> 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<String> = BUILD_TIMESTAMP.to_string().parse::<DateTime<Utc>>().ok().map(|datetime| datetime.format("%Y-%m-%d %H:%M:%S %Z").to_string());
Healthcheck {
+110 -77
View File
@@ -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<u16>,
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<String, ProviderLineup>,
providers: Vec<ProviderLineup>,
}
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<HashMap<String, u16>> {
let mut result = HashMap::<String, u16>::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<HashMap<String, u16>> {
let result = RefCell::new(HashMap::<String, u16>::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)
}
}
}
+2 -2
View File
@@ -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);
}
}
};
+1
View File
@@ -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;
@@ -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<AtomicOnceFlag>,
url: Url,
reconnect: bool,
reconnect_force_secs: u32,
connect_timeout_secs: u32,
headers: HeaderMap,
range_bytes: Arc<Option<AtomicUsize>>,
}
@@ -170,7 +179,7 @@ fn get_request_range_start_bytes(req_headers: &HashMap<String, Vec<u8>>) -> Opti
fn get_client_stream_request_params(
req_headers: &HeaderMap,
input: Option<&ConfigInput>,
options: &BufferStreamOptions) -> (usize, Option<usize>, bool, HeaderMap)
options: &BufferStreamOptions) -> (usize, Option<usize>, 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<reqwest::Client>, url: &Url, headers: &HeaderMap, range_start_bytes_to_request: Option<usize>) -> (reqwest::RequestBuilder, bool) {
@@ -248,15 +257,25 @@ async fn stream_provider(client: Arc<reqwest::Client>, 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,
}
@@ -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<Bytes, StreamError>;
fn poll_next(mut self: Pin<&mut Self>,cx: &mut std::task::Context<'_>,) -> Poll<Option<Self::Item>> {
if self.start_time.elapsed() > self.duration {
return Poll::Ready(None);
}
Pin::as_mut(&mut self.inner).poll_next(cx)
}
}
+5
View File
@@ -1108,6 +1108,11 @@ pub struct StreamConfig {
pub retry: bool,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub buffer: Option<StreamBufferConfig>,
#[serde(default)]
pub forced_retry_interval_secs: u32,
#[serde(default)]
pub connect_timeout_secs: u32,
}
impl StreamConfig {