diff --git a/.gitignore b/.gitignore index eef4fafd0..c0d5e4892 100644 --- a/.gitignore +++ b/.gitignore @@ -1,3 +1,4 @@ +/resources/*.ts /release /target /frontend/build diff --git a/CHANGELOG.md b/CHANGELOG.md index a10cf2e68..0cff6e7a9 100644 --- a/CHANGELOG.md +++ b/CHANGELOG.md @@ -1,5 +1,5 @@ # Changelog -# 2.1.4 (2025-01-xx) +# 2.2.0 (2025-02-11) - !BREAKING CHANGE! unique `input` `name` is now mandatory, because rearranging the `source.yml` could lead to wrong results without a playlist update. - !BREAKING_CHANGE! `log_sanitize_sensitive_info` is now under `log` section as `sanitize_sensitive_info` - !BREAKING_CHANGE! uuid generation for entries changed to `input.name` + `stream_id`. Virtual id mapping changed. The new Virtual id is not a sequence anymore. @@ -46,6 +46,12 @@ web_ui_enabled: true - added optional user properties: `max_connections`, `status`, `exp_date` (expiration date as unix seconds). If they exist they are checked when `config.yml` `user_access_control` set to true., if you don't need them remove this fields from `api-proxy.yml` Added option in `config.yml` the option `user_access_control` to activate the checks. Default is false. +- Added option `channel_unavailable_file` in `config.yml`. If a provider stream is not available this file content is send instead. +```yaml +update_on_boot: false +web_ui_enabled: true +channel_unavailable_file: /freeze_frame.ts +``` # 2.1.3 (2025-01-26) - Hotfix 2.1.2, forgot to update the stream api code. diff --git a/README.md b/README.md index ab8b1ccca..d85de662c 100644 --- a/README.md +++ b/README.md @@ -78,6 +78,7 @@ Top level entries in the config files are: * `reverse_proxy` _optional_ * `log` _optional * `user_access_control` _optional_ +* `channel_unavailable_file` _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. @@ -307,6 +308,20 @@ If you set it to `true`, the attributes (if available) are checked to permit or deny access. +### 1.12 `channel_unavailable_file` +If you want to send a `Unavailable Channel` picture instead of black screen when a channel is not available. +A video file with name `freeze_frame.ts` is already available in the docker image. + +You can convert an image with `ffmpeg`. + +`ffmpeg -loop 1 -framerate 1 -i freeze_frame.jpg -t 1 -c:v mpeg2video -f mpegts freeze_frame.ts` + +and add it to the `config.yml`. + +```yaml +channel_unavailable_file: /freeze_frame.ts +``` + ## 2. `source.yml` Has the following top level entries: diff --git a/bin/build_github_docker.sh b/bin/build_github_docker.sh index d321d1715..111b00003 100755 --- a/bin/build_github_docker.sh +++ b/bin/build_github_docker.sh @@ -3,6 +3,8 @@ set -euo pipefail source "${HOME}/.ghcr.io" WORKING_DIR=$(pwd) +BIN_DIR="${WORKING_DIR}/bin" +RESOURCES_DIR="${WORKING_DIR}/resources" DOCKER_DIR="${WORKING_DIR}/docker" FRONTEND_DIR="${WORKING_DIR}/frontend" TARGET=x86_64-unknown-linux-musl @@ -10,6 +12,11 @@ TARGET=x86_64-unknown-linux-musl cd "$FRONTEND_DIR" && rm -rf build && yarn && yarn build cd "$WORKING_DIR" +if [ ! -f "${BIN_DIR}/build_resources.sh" ]; then + "${BIN_DIR}/build_resources.sh" +fi + + # Check if the frontend build directory exists if [ ! -d "$FRONTEND_DIR/build" ]; then echo "Error: Web directory '$FRONTEND_DIR/build' does not exist." @@ -30,6 +37,7 @@ BIN_FILE=${WORKING_DIR}/target/${TARGET}/release/m3u-filter cp "${WORKING_DIR}/target/${TARGET}/release/m3u-filter" "${DOCKER_DIR}/" rm -rf "${DOCKER_DIR}/web" cp -r "${FRONTEND_DIR}/build" "${DOCKER_DIR}/web" +cp -r "${RESOURCES_DIR}/freeze_frame.ts" "${DOCKER_DIR}/" # Get the version from the binary VERSION=$("$BIN_FILE" -V | sed 's/m3u-filter *//') @@ -66,5 +74,6 @@ docker push ghcr.io/euzu/${ALPINE_IMAGE_NAME}:latest echo "Cleaning up build artifacts..." rm -rf "${DOCKER_DIR}/web" rm -f "${DOCKER_DIR}/m3u-filter" +rm -f "${DOCKER_DIR}/freeze_frame.ts" echo "Docker images for version ${VERSION} have been successfully built, tagged, and pushed." diff --git a/bin/build_github_docker_aarch64.sh b/bin/build_github_docker_aarch64.sh index e0cb876c6..ebef40283 100755 --- a/bin/build_github_docker_aarch64.sh +++ b/bin/build_github_docker_aarch64.sh @@ -3,6 +3,8 @@ set -euo pipefail source "${HOME}/.ghcr.io" WORKING_DIR=$(pwd) +BIN_DIR="${WORKING_DIR}/bin" +RESOURCES_DIR="${WORKING_DIR}/resources" DOCKER_DIR="${WORKING_DIR}/docker" FRONTEND_DIR="${WORKING_DIR}/frontend" TARGET=aarch64-unknown-linux-musl @@ -12,6 +14,11 @@ if [ -z "${VERSION}" ]; then exit 1 fi +if [ ! -f "${BIN_DIR}/build_resources.sh" ]; then + "${BIN_DIR}/build_resources.sh" +fi + + cd "$FRONTEND_DIR" && rm -rf build && yarn && yarn build cd "$WORKING_DIR" @@ -35,6 +42,7 @@ fi cp "${WORKING_DIR}/target/${TARGET}/release/m3u-filter" "${DOCKER_DIR}/" rm -rf "${DOCKER_DIR}/web" cp -r "${WORKING_DIR}/frontend/build" "${DOCKER_DIR}/web" +cp -r "${RESOURCES_DIR}/freeze_frame.ts" "${DOCKER_DIR}/" # Get the version from the binary # VERSION=$(./m3u-filter -V | sed 's/m3u-filter *//') @@ -58,5 +66,6 @@ docker push ghcr.io/euzu/${SCRATCH_IMAGE_NAME}:latest echo "Cleaning up build artifacts..." rm -rf "${DOCKER_DIR}/web" rm -f "${DOCKER_DIR}/m3u-filter" +rm -f "${DOCKER_DIR}/freeze_frame.ts" echo "Docker images ghcr.io/euzu/${SCRATCH_IMAGE_NAME}${VERSION} have been successfully built, tagged, and pushed." diff --git a/bin/build_github_docker_beta.sh b/bin/build_github_docker_beta.sh index 4373b094b..a17f2b879 100755 --- a/bin/build_github_docker_beta.sh +++ b/bin/build_github_docker_beta.sh @@ -5,7 +5,7 @@ source "${HOME}/.ghcr.io" TARGET=x86_64-unknown-linux-musl -(bin/build_fe.sh && bin/build_lin_static.sh) || exit 1; +(bin/build_fe.sh && bin/build_lin_static.sh && bin/build_resources.sh ) || exit 1; # Check if the binary exists if [ ! -f "./target/${TARGET}/release/m3u-filter" ]; then @@ -19,11 +19,17 @@ if [ ! -d "./frontend/build" ]; then exit 1 fi +if [ ! -f "./resources/freeze_frame.ts" ]; then + echo "Error: ./resources/freeze_frame.ts does not exist." + exit 1 +fi + # Prepare Docker build context cd ./docker cp "../target/${TARGET}/release/m3u-filter" . rm -rf ./web cp -r ../frontend/build ./web +cp -r ../resources/freeze_frame.ts . # Get the version from the binary VERSION=$(./m3u-filter -V | sed 's/m3u-filter *//') @@ -58,5 +64,6 @@ docker push ghcr.io/euzu/m3u-filter-beta:latest echo "Cleaning up build artifacts..." rm -rf ./web rm -f ./m3u-filter +rm -f ./freeze_frame.ts echo "Docker images for version ${VERSION} have been successfully built, tagged, and pushed." diff --git a/bin/build_resources.sh b/bin/build_resources.sh new file mode 100755 index 000000000..493ab8985 --- /dev/null +++ b/bin/build_resources.sh @@ -0,0 +1,14 @@ +#!/usr/bin/env bash + +if [ -e ./resources/freeze_frame.ts ]; then + echo "Resource exists, skipping creation" + exit; +fi + + +if which ffmpeg > /dev/null 2>&1; then + ffmpeg -loop 1 -framerate 1 -i ./resources/freeze_frame.jpg -t 1 -c:v mpeg2video -f mpegts ./resources/freeze_frame.ts +else + echo "ffmpeg not found"; + exit; +fi diff --git a/config/config.yml b/config/config.yml index cb953a004..1a9a2dff0 100644 --- a/config/config.yml +++ b/config/config.yml @@ -3,6 +3,7 @@ threads: 0 working_dir: /home/m3u-filter/data backup_dir: /home/m3u-filter/.backup update_on_boot: false +channel_unavailable_file: /home/m3u-filter/freeze_frame.ts # sec min hour day of month month day of week year schedules: - schedule: "0 0 8,12,16,20,22,1 * * * *" diff --git a/docker/Dockerfile b/docker/Dockerfile index e3032e70a..3322f9a5c 100644 --- a/docker/Dockerfile +++ b/docker/Dockerfile @@ -39,6 +39,7 @@ COPY --from=rust-build /etc/ssl/certs/ca-certificates.crt /etc/ssl/certs/ COPY --from=rust-build /src/target/x86_64-unknown-linux-musl/release/m3u-filter /m3u-filter COPY --from=node-build /app/build /web +COPY ./freeze_frame.ts /freeze_frame.ts # config should be mounted as volume # COPY ./config /config @@ -54,6 +55,7 @@ WORKDIR /app COPY --from=rust-build /src/target/x86_64-unknown-linux-musl/release/m3u-filter m3u-filter COPY --from=node-build /app/build web +COPY ./freeze_frame.ts /freeze_frame.ts # config should be mounted as volume # COPY ./config config diff --git a/docker/Dockerfile-manual b/docker/Dockerfile-manual index 92b0f322b..a4730964e 100644 --- a/docker/Dockerfile-manual +++ b/docker/Dockerfile-manual @@ -9,6 +9,7 @@ COPY --from=build /etc/ssl/certs/ca-certificates.crt /etc/ssl/certs/ COPY ./m3u-filter / COPY ./web /web +COPY ./freeze_frame.ts ./freeze_frame.ts CMD ["./m3u-filter", "-s", "-p", "/config"] @@ -21,6 +22,7 @@ WORKDIR /app COPY ./m3u-filter . COPY ./web ./web +COPY ./freeze_frame.ts ./freeze_frame.ts # config should be mounted as volume # COPY ./config ./config diff --git a/resources/freeze_frame.jpg b/resources/freeze_frame.jpg new file mode 100644 index 000000000..c79ae5bba Binary files /dev/null and b/resources/freeze_frame.jpg differ diff --git a/src/api/api_utils.rs b/src/api/api_utils.rs index 6e46d5faf..2655e63a9 100644 --- a/src/api/api_utils.rs +++ b/src/api/api_utils.rs @@ -2,7 +2,6 @@ use crate::api::model::app_state::AppState; use crate::api::model::model_utils::get_stream_response_with_headers; use crate::api::model::streams::persist_pipe_stream::PersistPipeStream; use crate::api::model::streams::provider_stream; -use crate::api::model::streams::provider_stream::get_provider_pipe_stream; use crate::api::model::streams::provider_stream_factory::BufferStreamOptions; use crate::api::model::request::UserApiRequest; use crate::api::model::stream_error::StreamError; @@ -10,7 +9,7 @@ use crate::utils::{debug_if_enabled, trace_if_enabled}; use crate::model::api_proxy::ProxyUserCredentials; use crate::model::config::{ConfigInput, ConfigTarget}; use crate::model::playlist::PlaylistItemType; -use crate::utils::file::file_utils::create_new_file_for_write; +use crate::utils::file::file_utils::{create_new_file_for_write}; use crate::tools::lru_cache::LRUResourceCache; use crate::utils::network::request; use crate::utils::network::request::sanitize_sensitive_info; @@ -23,6 +22,7 @@ use futures::TryStreamExt; use log::{error, log_enabled, trace}; use reqwest::StatusCode; use std::collections::HashMap; +use std::io::BufWriter; use std::path::Path; use std::sync::Arc; use url::Url; @@ -103,7 +103,7 @@ 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, bool) { +fn get_stream_options(app_state: &AppState) -> (bool, bool, usize, bool) { let (stream_retry, buffer_enabled, buffer_size) = app_state .config .reverse_proxy @@ -117,8 +117,7 @@ fn get_stream_options(app_state: &AppState) -> (bool, bool, usize, bool, bool) { (stream.retry, buffer_enabled, buffer_size) }); let pipe_provider_stream = !stream_retry && !buffer_enabled; - let shared_stream_use_own_buffer = !buffer_enabled || pipe_provider_stream; - (stream_retry, buffer_enabled, buffer_size, pipe_provider_stream, shared_stream_use_own_buffer) + (stream_retry, buffer_enabled, buffer_size, pipe_provider_stream) } fn get_stream_content_length(provider_response: Option<&(Vec<(String, String)>, StatusCode)>) -> u64 { @@ -144,23 +143,23 @@ pub async fn stream_response(app_state: &AppState, stream_url: &str, } } - let (stream_retry, buffer_enabled, buffer_size, direct_pipe_provider_stream, shared_stream_use_own_buffer) = + let (stream_retry, buffer_enabled, buffer_size, direct_pipe_provider_stream) = get_stream_options(app_state); if let Ok(url) = Url::parse(stream_url) { let active_clients = Arc::clone(&app_state.active_users); let (stream_opt, provider_response) = if direct_pipe_provider_stream { - get_provider_pipe_stream(&app_state.http_client, &url, req, input, item_type).await + provider_stream::get_provider_pipe_stream(&app_state.config, &app_state.http_client, &url, req, input, item_type).await } else { - let buffer_stream_options = BufferStreamOptions::new(item_type, stream_retry, buffer_enabled, buffer_size); - provider_stream::get_provider_reconnect_buffered_stream(&app_state.http_client, &url, req, input, buffer_stream_options).await + let buffer_stream_options = BufferStreamOptions::new(item_type, stream_retry, buffer_enabled, buffer_size, share_stream); + provider_stream::get_provider_reconnect_buffered_stream(&app_state.config, &app_state.http_client, &url, req, input, buffer_stream_options).await }; if let Some(stream) = stream_opt { let content_length = get_stream_content_length(provider_response.as_ref()); let stream = ActiveClientStream::new(stream, active_clients, user, log_active_clients); let stream_resp = if share_stream { let shared_headers = provider_response.as_ref().map_or_else(Vec::new, |(h, _)| h.clone()); - SharedStreamManager::subscribe(app_state, stream_url, stream, shared_stream_use_own_buffer, shared_headers); + SharedStreamManager::subscribe(app_state, stream_url, stream, shared_headers, buffer_size); if let Some(broadcast_stream) = SharedStreamManager::subscribe_shared_stream(app_state, stream_url) { let mut response_builder = get_stream_response_with_headers(provider_response, stream_url); if content_length > 0 { @@ -261,7 +260,7 @@ pub async fn resource_response(app_state: &AppState, resource_url: &str, req: &H cache.lock().store_path(resource_url) }; if let Ok(file) = create_new_file_for_write(&resource_path) { - let writer = Arc::new(file); + let writer = BufWriter::new(file); let add_cache_content = get_add_cache_content(resource_url, &app_state.cache); let stream = PersistPipeStream::new(byte_stream, writer, add_cache_content); return response_builder.body(BodyStream::new(stream)); diff --git a/src/api/model/active_user_manager.rs b/src/api/model/active_user_manager.rs index ba41702a3..c8b4cd91a 100644 --- a/src/api/model/active_user_manager.rs +++ b/src/api/model/active_user_manager.rs @@ -42,6 +42,7 @@ impl ActiveUserManager { } else { lock.insert(username.to_string(), AtomicU32::new(1)); } + drop(lock); } (self.active_users(), self.active_connections()) } @@ -54,6 +55,7 @@ impl ActiveUserManager { lock.remove(username); } } + drop(lock); } (self.active_users(), self.active_connections()) } diff --git a/src/api/model/streams/active_client_stream.rs b/src/api/model/streams/active_client_stream.rs index c7ffbff46..b75227318 100644 --- a/src/api/model/streams/active_client_stream.rs +++ b/src/api/model/streams/active_client_stream.rs @@ -18,9 +18,7 @@ pub(in crate::api) struct ActiveClientStream { impl ActiveClientStream { pub(crate) fn new(inner: ResponseStream, active_clients: Arc, user: &ProxyUserCredentials, log_active_clients: bool) -> Self { - let (client_count, connection_count) = { - active_clients.add_connection(&user.username) - }; + let (client_count, connection_count) = active_clients.add_connection(&user.username); if log_active_clients { info!("Active clients: {client_count}, active connections {connection_count}"); } @@ -30,10 +28,7 @@ impl ActiveClientStream { impl Stream for ActiveClientStream { type Item = Result; - fn poll_next( - mut self: Pin<&mut Self>, - cx: &mut std::task::Context<'_>, - ) -> Poll> { + fn poll_next(mut self: Pin<&mut Self>,cx: &mut std::task::Context<'_>,) -> Poll> { Pin::as_mut(&mut self.inner).poll_next(cx) } } @@ -41,9 +36,7 @@ impl Stream for ActiveClientStream { impl Drop for ActiveClientStream { fn drop(&mut self) { - let (client_count, connection_count) = { - self.active_clients.remove_connection(&self.username) - }; + let (client_count, connection_count) = self.active_clients.remove_connection(&self.username); if self.log_active_clients { info!("Active clients: {client_count}, active connections {connection_count}"); } diff --git a/src/api/model/streams/buffered_stream.rs b/src/api/model/streams/buffered_stream.rs index 7c358a16b..2bc98aad7 100644 --- a/src/api/model/streams/buffered_stream.rs +++ b/src/api/model/streams/buffered_stream.rs @@ -4,8 +4,7 @@ use std::{ pin::Pin, sync::Arc, }; -use std::time::Duration; -use tokio::sync::mpsc::channel; +use tokio::sync::mpsc::{channel, Sender}; use tokio_stream::wrappers::ReceiverStream; use crate::api::model::stream_error::StreamError; use crate::tools::atomic_once_flag::AtomicOnceFlag; @@ -17,33 +16,43 @@ pub(in crate::api::model) struct BufferedStream { impl BufferedStream { pub fn new(stream: ResponseStream, buffer_size: usize, client_close_signal: Arc, _url: &str) -> Self { let (tx, rx) = channel(buffer_size); - actix_rt::spawn(async move { - let mut stream = stream; - let sleep_duration= Duration::from_millis(100); - loop { - match stream.next().await { - Some(Ok(chunk)) => { - // this is for backpressure, we fill the buffer and wait for the receiver - if let Ok(permit) = tx.reserve().await { - permit.send(Ok(chunk)); - } else { - // receiver closed. + actix_rt::spawn(Self::buffer_stream(tx, stream, client_close_signal)); + Self { + stream: ReceiverStream::new(rx) + } + } + + async fn buffer_stream( + tx: Sender>, + mut stream: ResponseStream, + client_close_signal: Arc, + ) { + loop { + if !client_close_signal.is_active() { + break; + } + match stream.next().await { + Some(Ok(chunk)) => { + match tx.reserve().await { + Ok(permit) => permit.send(Ok(chunk)), + Err(_err) => { + // Receiver dropped, notify and exit client_close_signal.notify(); break; } } - Some(Err(_err)) => { - actix_web::rt::time::sleep(sleep_duration).await; - } - None => { - break - } } + Some(Err(err)) => { + eprintln!("Buffered Stream Error: {err:?}"); + // actix_web::rt::time::sleep(sleep_duration).await; + // Attempt to send error to client + if tx.send(Err(err)).await.is_err() { + client_close_signal.notify(); + } + break; + } + None => break, } - }); - - Self { - stream: ReceiverStream::new(rx) } } } @@ -51,7 +60,7 @@ impl BufferedStream { impl Stream for BufferedStream { type Item = Result; - fn poll_next(mut self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll> { - self.stream.poll_next_unpin(cx) + fn poll_next(self: Pin<&mut Self>, cx: &mut Context<'_>) -> Poll> { + Pin::new(&mut self.get_mut().stream).poll_next(cx) } } diff --git a/src/api/model/streams/client_stream.rs b/src/api/model/streams/client_stream.rs index 4a80b912d..eb8bc6ef4 100644 --- a/src/api/model/streams/client_stream.rs +++ b/src/api/model/streams/client_stream.rs @@ -35,6 +35,7 @@ impl Stream for ClientStream { match Pin::as_mut(&mut self.inner).poll_next(cx) { Poll::Ready(Some(Ok(bytes))) => { if bytes.is_empty() { + eprintln!("client stream empty bytes"); continue; } @@ -48,7 +49,10 @@ impl Stream for ClientStream { self.close_signal.notify(); return Poll::Ready(None); } - other => return other, + Poll::Pending => return Poll::Pending, + Poll::Ready(Some(Err(err))) => { + eprintln!("client stream error: {err}"); + } } } } diff --git a/src/api/model/streams/freeze_frame_stream.rs b/src/api/model/streams/freeze_frame_stream.rs new file mode 100644 index 000000000..2e09ca857 --- /dev/null +++ b/src/api/model/streams/freeze_frame_stream.rs @@ -0,0 +1,57 @@ +use crate::api::model::stream_error::StreamError; +use bytes::Bytes; +use futures::Stream; +use std::pin::Pin; +use std::sync::Arc; +use std::task::{Context, Poll}; + +const CHUNK_SIZE: usize = 8192; + + +pub struct FreezeFrameStream { + buffer: Arc>, + buffer_len: usize, + current_pos: usize, // Keep track of the current position in the buffer +} + +impl FreezeFrameStream { + pub fn new(_status: u16, buffer: Arc>) -> Self { + let buffer_len = buffer.len(); + Self { + buffer, + buffer_len, + current_pos: 0, + } + } +} + +impl Stream for FreezeFrameStream { + type Item = Result; + + fn poll_next( + mut self: Pin<&mut Self>, + _cx: &mut Context<'_>, + ) -> Poll> { + + if self.buffer_len == 0 { + return Poll::Ready(None); // If buffer is empty, return None (end of stream) + } + + // Calculate the start and end positions for the chunk + let start = self.current_pos; + let end = (self.current_pos + CHUNK_SIZE).min(self.buffer_len); + + // Create a chunk from the buffer + let chunk = self.buffer[start..end].to_vec(); + let bytes_chunk = Bytes::from(chunk); + + // Update the current position + self.current_pos = if end == self.buffer_len { + 0 // Wrap around if we reach the end + } else { + end + }; + + Poll::Ready(Some(Ok(bytes_chunk))) + } +} diff --git a/src/api/model/streams/mod.rs b/src/api/model/streams/mod.rs index 42f8c712e..e290ed13d 100644 --- a/src/api/model/streams/mod.rs +++ b/src/api/model/streams/mod.rs @@ -4,4 +4,5 @@ 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 buffered_stream; -mod client_stream; \ No newline at end of file +mod client_stream; +mod freeze_frame_stream; \ No newline at end of file diff --git a/src/api/model/streams/persist_pipe_stream.rs b/src/api/model/streams/persist_pipe_stream.rs index bba910a95..525c1b9a7 100644 --- a/src/api/model/streams/persist_pipe_stream.rs +++ b/src/api/model/streams/persist_pipe_stream.rs @@ -91,4 +91,4 @@ where } } } -} \ No newline at end of file +} diff --git a/src/api/model/streams/provider_stream.rs b/src/api/model/streams/provider_stream.rs index b6eb9889a..d9be0b370 100644 --- a/src/api/model/streams/provider_stream.rs +++ b/src/api/model/streams/provider_stream.rs @@ -1,19 +1,20 @@ use crate::api::api_utils::{get_headers_from_request, HeaderFilter}; +use crate::api::model::model_utils::get_response_headers; +use crate::api::model::stream_error::StreamError; +use crate::api::model::streams::freeze_frame_stream::FreezeFrameStream; use crate::api::model::streams::provider_stream_factory::{create_provider_stream, BufferStreamOptions}; -use crate::model::config::ConfigInput; +use crate::model::config::{Config, ConfigInput}; +use crate::model::playlist::PlaylistItemType; +use crate::utils::debug_if_enabled; use crate::utils::network::request::{get_request_headers, sanitize_sensitive_info}; -use actix_web::{HttpRequest}; +use actix_web::HttpRequest; use bytes::Bytes; use futures::stream::BoxStream; +use futures::TryStreamExt; use log::error; use reqwest::StatusCode; use std::sync::Arc; -use futures::TryStreamExt; use url::Url; -use crate::api::model::model_utils::get_response_headers; -use crate::api::model::stream_error::StreamError; -use crate::model::playlist::PlaylistItemType; -use crate::utils::debug_if_enabled; type ProviderStreamResponse = (Option>>, Option<(Vec<(String, String)>, StatusCode)>); @@ -27,7 +28,8 @@ pub fn get_header_filter_for_item_type(item_type: PlaylistItemType) -> HeaderFil } } -pub async fn get_provider_pipe_stream(http_client: &Arc, +pub async fn get_provider_pipe_stream(cfg: &Config, + http_client: &Arc, stream_url: &Url, req: &HttpRequest, input: Option<&ConfigInput>, @@ -46,7 +48,9 @@ pub async fn get_provider_pipe_stream(http_client: &Arc, let response_headers = get_response_headers(&mut response); let status = response.status(); if status.is_success() { - (Some(Box::pin(response.bytes_stream().map_err(|err|StreamError::reqwest(&err)))), Some((response_headers, status))) + (Some(Box::pin(response.bytes_stream().map_err(|err| StreamError::reqwest(&err)))), Some((response_headers, status))) + } else if let Some(freeze_frame) = cfg.t_channel_unavailable_file.as_ref() { + (Some(Box::pin(FreezeFrameStream::new(status.as_u16(), Arc::clone(freeze_frame)))), Some((response_headers, status))) } else { (None, Some((response_headers, status))) } @@ -59,12 +63,13 @@ pub async fn get_provider_pipe_stream(http_client: &Arc, } } -pub async fn get_provider_reconnect_buffered_stream(http_client: &Arc, +pub async fn get_provider_reconnect_buffered_stream(cfg: &Config, + http_client: &Arc, stream_url: &Url, req: &HttpRequest, input: Option<&ConfigInput>, options: BufferStreamOptions) -> ProviderStreamResponse { - match create_provider_stream(Arc::clone(http_client), stream_url, req, input, options).await { + match create_provider_stream(cfg, Arc::clone(http_client), stream_url, req, input, options).await { None => (None, None), Some((stream, info)) => { (Some(stream), info) diff --git a/src/api/model/streams/provider_stream_factory.rs b/src/api/model/streams/provider_stream_factory.rs index 31f8bd965..88649442d 100644 --- a/src/api/model/streams/provider_stream_factory.rs +++ b/src/api/model/streams/provider_stream_factory.rs @@ -1,12 +1,13 @@ use crate::api::api_utils::get_headers_from_request; -use crate::api::model::streams::buffered_stream::BufferedStream; -use crate::api::model::streams::client_stream::ClientStream; use crate::api::model::model_utils::get_response_headers; use crate::api::model::stream_error::StreamError; -use crate::utils::debug_if_enabled; -use crate::model::config::ConfigInput; +use crate::api::model::streams::buffered_stream::BufferedStream; +use crate::api::model::streams::client_stream::ClientStream; +use crate::api::model::streams::provider_stream::get_header_filter_for_item_type; +use crate::model::config::{Config, ConfigInput}; use crate::model::playlist::PlaylistItemType; use crate::tools::atomic_once_flag::AtomicOnceFlag; +use crate::utils::debug_if_enabled; use crate::utils::network::request::{classify_content_type, get_request_headers, sanitize_sensitive_info, MimeCategory}; use actix_web::HttpRequest; use bytes::Bytes; @@ -20,7 +21,7 @@ use std::sync::atomic::{AtomicUsize, Ordering}; use std::sync::Arc; use std::time::{Duration, Instant}; use url::Url; -use crate::api::model::streams::provider_stream::get_header_filter_for_item_type; +use crate::api::model::streams::freeze_frame_stream::FreezeFrameStream; // TODO make this configurable pub const STREAM_QUEUE_SIZE: usize = 1024; // mpsc channel holding messages. with 8092byte chunks and 2Mbit/s approx 8MB @@ -35,6 +36,7 @@ pub struct BufferStreamOptions { reconnect_enabled: bool, buffer_enabled: bool, buffer_size: usize, + share_stream: bool, } impl BufferStreamOptions { @@ -43,12 +45,14 @@ impl BufferStreamOptions { reconnect_enabled: bool, buffer_enabled: bool, buffer_size: usize, + share_stream: bool ) -> Self { Self { item_type, reconnect_enabled, buffer_enabled, buffer_size, + share_stream, } } @@ -57,6 +61,11 @@ impl BufferStreamOptions { self.buffer_enabled } + #[inline] + fn is_shared_stream(&self) -> bool { + self.share_stream + } + // #[inline] // fn get_buffer_size(&self) -> usize { // self.buffer_size @@ -193,7 +202,7 @@ fn prepare_client(request_client: &Arc, url: &Url, headers: &He } } -async fn provider_request(request_client: Arc, initial_info: bool, stream_options: &ProviderStreamOptions) -> Result, StatusCode> { +async fn provider_request(cfg: &Config, request_client: Arc, initial_info: bool, stream_options: &ProviderStreamOptions) -> Result, StatusCode> { let (client, _partial_content) = prepare_client(&request_client, stream_options.get_url(), stream_options.get_headers(), stream_options.get_total_bytes_send()); match client.send().await { Ok(mut response) => { @@ -214,6 +223,9 @@ async fn provider_request(request_client: Arc, initial_info: bo StreamError::reqwest(&err) }).boxed(), response_info))); } + if let Some(freeze_frame) = cfg.t_channel_unavailable_file.as_ref() { + return Ok(Some((FreezeFrameStream::new(status.as_u16(), Arc::clone(freeze_frame)).boxed(), None))) + } Err(status) } Err(_err) => { @@ -260,17 +272,17 @@ async fn stream_provider(client: Arc, stream_options: ProviderS } actix_web::rt::time::sleep(Duration::from_millis(100)).await; } - debug_if_enabled!("Stopped seconnecting stream {}", sanitize_sensitive_info(url.as_str())); + debug_if_enabled!("Stopped reconnecting stream {}", sanitize_sensitive_info(url.as_str())); None } const RETRY_SECONDS: u64 = 5; const ERR_MAX_RETRY_COUNT: u32 = 5; -async fn get_initial_stream(client: Arc, stream_options: &ProviderStreamOptions) -> Option { +async fn get_initial_stream(cfg: &Config, client: Arc, stream_options: &ProviderStreamOptions) -> Option { let start = Instant::now(); let mut connect_err: u32 = 1; while stream_options.should_continue() { - match provider_request(Arc::clone(&client), true, stream_options).await { + match provider_request(cfg, Arc::clone(&client), true, stream_options).await { Ok(Some(value)) => return Some(value), Ok(None) => { if connect_err > ERR_MAX_RETRY_COUNT { @@ -316,7 +328,8 @@ fn create_provider_stream_options(stream_url: &Url, } } -pub async fn create_provider_stream(client: Arc, +pub async fn create_provider_stream(cfg: &Config, + client: Arc, stream_url: &Url, req: &HttpRequest, input: Option<&ConfigInput>, @@ -324,7 +337,7 @@ pub async fn create_provider_stream(client: Arc, let stream_options = create_provider_stream_options(stream_url, req, input, &options); let client_stream_factory = |stream, reconnect_flag, range_cnt| { - let stream = if stream_options.is_buffered() { + let stream = if stream_options.is_buffered() && !options.is_shared_stream() { BufferedStream::new(stream, stream_options.get_buffer_size(), stream_options.get_continue_flag_clone(), stream_url.as_str()).boxed() } else { stream @@ -332,7 +345,7 @@ pub async fn create_provider_stream(client: Arc, ClientStream::new(stream, reconnect_flag, range_cnt, stream_options.get_url().as_str()).boxed() }; - match get_initial_stream(Arc::clone(&client), &stream_options).await { + match get_initial_stream(cfg, Arc::clone(&client), &stream_options).await { Some((init_stream, info)) => { let is_media_stream = if let Some((headers, _)) = &info { classify_content_type(headers) == MimeCategory::Video @@ -344,12 +357,18 @@ pub async fn create_provider_stream(client: Arc, if is_media_stream && stream_options.should_reconnect() { let client_signal = Arc::clone(&continue_signal); let stream_options_provider = stream_options.clone(); + let continue_streaming_signal = client_signal.clone(); let unfold: ResponseStream = stream::unfold((), move |()| { let client = Arc::clone(&client); let stream_opts = stream_options_provider.clone(); + let continue_streaming = continue_streaming_signal.clone(); async move { - let stream = stream_provider(client, stream_opts).await?; - Some((stream, ())) + if continue_streaming.is_active() { + let stream = stream_provider(client, stream_opts).await?; + Some((stream, ())) + } else { + None + } } }).flatten().boxed(); Some((client_stream_factory(init_stream.chain(unfold).boxed(), Arc::clone(&client_signal), stream_options.get_range_bytes_clone()).boxed(), info)) diff --git a/src/api/model/streams/shared_stream_manager.rs b/src/api/model/streams/shared_stream_manager.rs index 5234f0ae1..3f6dcae25 100644 --- a/src/api/model/streams/shared_stream_manager.rs +++ b/src/api/model/streams/shared_stream_manager.rs @@ -3,7 +3,7 @@ use crate::api::model::streams::provider_stream_factory::STREAM_QUEUE_SIZE; use crate::api::model::stream_error::StreamError; use crate::utils::debug_if_enabled; use crate::utils::network::request::sanitize_sensitive_info; -use parking_lot::{FairMutex}; +use parking_lot::{RwLock}; use bytes::Bytes; use futures::stream::BoxStream; use futures::{Stream, StreamExt}; @@ -15,8 +15,6 @@ use std::pin::Pin; use std::task::{Context, Poll}; use std::time::Duration; -const MIN_STREAM_QUEUE_SIZE: usize = 128; - /// /// Wraps a `ReceiverStream` as Stream> /// @@ -67,6 +65,7 @@ impl SharedStreamState { fn broadcast(&self, stream_url: &str, bytes_stream: S, shared_streams: Arc) where S: Stream> + Unpin + 'static, + E: std::fmt::Debug { let mut source_stream = Box::pin(bytes_stream); let streaming_url = stream_url.to_string(); @@ -81,9 +80,13 @@ impl SharedStreamState { debug_if_enabled!("No active subscribers. Closing shared provider stream {}", sanitize_sensitive_info(&streaming_url)); break; } - let _ = sender.send(data); + if let Err(err) = sender.send(data) { + eprintln!("broadcast send err {err:?}"); + } } - None | Some(Err(_)) => { + None => break, + Some(Err(err)) => { + eprintln!("broadcast failure {err:?}"); break; } } @@ -94,7 +97,7 @@ impl SharedStreamState { } } -type SharedStreamRegister = Arc>>; +type SharedStreamRegister = RwLock>; pub struct SharedStreamManager { shared_streams: SharedStreamRegister, @@ -103,35 +106,36 @@ pub struct SharedStreamManager { impl SharedStreamManager { pub(crate) fn new() -> Self { Self { - shared_streams: Arc::new(FairMutex::new(HashMap::new())), + shared_streams: RwLock::new(HashMap::new()), } } pub fn get_shared_state_headers(&self, stream_url: &str) -> Option> { - self.shared_streams.lock().get(stream_url).map(|s| s.headers.clone()) + self.shared_streams.read().get(stream_url).map(|s| s.headers.clone()) } fn unregister(&self, stream_url: &str) { - self.shared_streams.lock().remove(stream_url); + self.shared_streams.write().remove(stream_url); + } + + fn register(&self, stream_url: &str, shared_state: SharedStreamState) { + self.shared_streams.write().insert(stream_url.to_string(), shared_state); } pub(crate) fn subscribe( app_state: &AppState, stream_url: &str, bytes_stream: S, - use_buffer: bool, - headers: Vec<(String, String)>, ) + headers: Vec<(String, String)>, + buffer_size: usize,) where S: Stream> + Unpin + 'static, + E: std::fmt::Debug { - let buf_size = if use_buffer { STREAM_QUEUE_SIZE } else { MIN_STREAM_QUEUE_SIZE }; + let buf_size = std::cmp::max(buffer_size, STREAM_QUEUE_SIZE); let shared_state = SharedStreamState::new(headers, buf_size); shared_state.broadcast(stream_url, bytes_stream, Arc::clone(&app_state.shared_stream_manager)); - app_state - .shared_stream_manager - .shared_streams - .lock() - .insert(stream_url.to_string(), shared_state); + app_state.shared_stream_manager.register(stream_url, shared_state); debug_if_enabled!("Created shared provider stream {}", sanitize_sensitive_info(stream_url)); } @@ -140,15 +144,7 @@ impl SharedStreamManager { app_state: &AppState, stream_url: &str, ) -> Option>> { - if let Some(shared_stream) = app_state - .shared_stream_manager - .shared_streams - .lock() - .get(stream_url) { - debug_if_enabled!("Responding existing shared client stream {}", sanitize_sensitive_info(stream_url)); - Some(shared_stream.subscribe()) - } else { - None - } + debug_if_enabled!("Responding existing shared client stream {}", sanitize_sensitive_info(stream_url)); + app_state.shared_stream_manager.shared_streams.read().get(stream_url).map(SharedStreamState::subscribe) } } \ No newline at end of file diff --git a/src/model/config.rs b/src/model/config.rs index b2470076d..47cb6ae1f 100644 --- a/src/model/config.rs +++ b/src/model/config.rs @@ -1007,6 +1007,8 @@ pub struct Config { #[serde(default, skip_serializing_if = "Option::is_none")] pub backup_dir: Option, #[serde(default, skip_serializing_if = "Option::is_none")] + pub channel_unavailable_file: Option, + #[serde(default, skip_serializing_if = "Option::is_none")] pub templates: Option>, #[serde(default, skip_serializing_if = "Option::is_none")] pub video: Option, @@ -1038,6 +1040,8 @@ pub struct Config { pub t_api_proxy_file_path: String, #[serde(skip)] pub file_locks: Arc, + #[serde(skip)] + pub t_channel_unavailable_file: Option>>, } impl Config { @@ -1172,6 +1176,17 @@ impl Config { let work_dir = if resolve_var { &config_reader::resolve_env_var(&self.working_dir) } else { &self.working_dir }; self.working_dir = file_utils::get_working_path(work_dir); + if let Some(channel_unavailable_file) = &self.channel_unavailable_file { + let channel_unavailable = file_utils::make_absolute_path(channel_unavailable_file, &self.working_dir, resolve_var); + match file_utils::read_file_as_bytes(&PathBuf::from(&channel_unavailable)) { + Ok(data) => self.t_channel_unavailable_file = Some(Arc::new(data)), + Err(err) => { + error!("Failed to load channel unavailable file: {channel_unavailable} {err}"); + } + } + self.channel_unavailable_file = Some(channel_unavailable); + } + if self.backup_dir.is_none() { self.backup_dir = Some(PathBuf::from(&self.working_dir).join("backup").clean().to_string_lossy().to_string()); } else { @@ -1250,24 +1265,7 @@ impl Config { fn prepare_api_web_root(&mut self, resolve_var: bool) { if !self.api.web_root.is_empty() { - let web_root = if resolve_var { config_reader::resolve_env_var(&self.api.web_root) } else { self.api.web_root.clone() }; - self.api.web_root = web_root.to_string(); - let wrpb = std::path::PathBuf::from(&self.api.web_root); - if wrpb.is_relative() { - let mut wrpb2 = std::path::PathBuf::from(&self.working_dir).join(&web_root); - if !wrpb2.exists() { - wrpb2 = file_utils::get_exe_path().join(&web_root); - } - if !wrpb2.exists() { - let cwd = std::env::current_dir(); - if let Ok(cwd_path) = cwd { - wrpb2 = cwd_path.join(&web_root); - } - } - if wrpb2.exists() { - self.api.web_root = String::from(wrpb2.clean().to_str().unwrap_or_default()); - } - } + self.api.web_root = file_utils::make_absolute_path(&self.api.web_root, &self.working_dir, resolve_var); } } diff --git a/src/tools/atomic_once_flag.rs b/src/tools/atomic_once_flag.rs index f13e36162..7066878ff 100644 --- a/src/tools/atomic_once_flag.rs +++ b/src/tools/atomic_once_flag.rs @@ -16,7 +16,6 @@ use std::sync::atomic::{AtomicBool, Ordering}; #[derive(Debug)] pub struct AtomicOnceFlag { enabled: AtomicBool, - ordering: Ordering, } impl Default for AtomicOnceFlag { @@ -26,30 +25,24 @@ impl Default for AtomicOnceFlag { } impl AtomicOnceFlag { - /// Creates a new `AtomicOnceFlag` with the specified memory ordering. - pub fn with_ordering(ordering: Ordering) -> Self { - Self { - enabled: AtomicBool::new(true), - ordering, - } - } - /// Creates a new `AtomicOnceFlag` with a default memory ordering of `Relaxed`. pub fn new() -> Self { - Self::with_ordering(Ordering::SeqCst) + Self { + enabled: AtomicBool::new(true), + } } /// Disables the flag. After calling this method, `is_active()` will always return `false`. /// /// This operation is atomic and uses the specified memory ordering. pub fn notify(&self) { - self.enabled.store(false, self.ordering); + self.enabled.store(false, Ordering::SeqCst); } /// Checks if the flag is still active. /// /// Returns `true` if the flag is active (initial state). Returns `false` if the flag has been disabled. pub fn is_active(&self) -> bool { - self.enabled.load(self.ordering) + self.enabled.load(Ordering::SeqCst) } } \ No newline at end of file diff --git a/src/utils/file/file_utils.rs b/src/utils/file/file_utils.rs index da6bb85c3..f85e73110 100644 --- a/src/utils/file/file_utils.rs +++ b/src/utils/file/file_utils.rs @@ -4,10 +4,11 @@ use std::fs::{File, OpenOptions}; use std::io::{BufReader, BufWriter, Read, Write}; use std::path::{Path, PathBuf}; +use crate::m3u_filter_error::str_to_io_error; +use crate::utils::debug_if_enabled; +use crate::utils::file::{config_reader}; use log::{debug, error}; use path_clean::PathClean; -use crate::utils::debug_if_enabled; -use crate::m3u_filter_error::str_to_io_error; const USER_FILE: &str = "user.txt"; const CONFIG_PATH: &str = "config"; @@ -17,13 +18,15 @@ const MAPPING_FILE: &str = "mapping.yml"; const API_PROXY_FILE: &str = "api-proxy.yml"; pub fn file_writer(w: W) -> BufWriter -where W: Write +where + W: Write, { BufWriter::with_capacity(131_072, w) } pub fn file_reader(r: R) -> BufReader -where R: Read +where + R: Read, { BufReader::with_capacity(131_072, r) } @@ -250,4 +253,33 @@ pub fn prepare_file_path(persist: Option<&str>, working_dir: &str, action: &str) } else { None } +} + +pub fn read_file_as_bytes(path: &Path) -> std::io::Result> { + let mut file = File::open(path)?; + let mut buffer = Vec::new(); + file.read_to_end(&mut buffer)?; + Ok(buffer) +} + +pub fn make_absolute_path(path: &str, working_dir: &str, resolve_var: bool) -> String { + let resolved_path = if resolve_var { config_reader::resolve_env_var(path) } else { path.to_string() }; + + let rpb = std::path::PathBuf::from(&resolved_path); + if rpb.is_relative() { + let mut rpb2 = std::path::PathBuf::from(working_dir).join(&rpb); + if !rpb2.exists() { + rpb2 = get_exe_path().join(&rpb); + } + if !rpb2.exists() { + let cwd = std::env::current_dir(); + if let Ok(cwd_path) = cwd { + rpb2 = cwd_path.join(&rpb); + } + } + if rpb2.exists() { + return String::from(rpb2.clean().to_str().unwrap_or_default()); + } + } + resolved_path } \ No newline at end of file