diff --git a/README.md b/README.md index 66fccb337..f12f81b7e 100644 --- a/README.md +++ b/README.md @@ -218,7 +218,7 @@ Attributes: - `throttle` Allowed units are `KB/s`,`MB/s`,`KiB/s`,`MiB/s`,`kbps`,`mbps`,`Mibps`. Default unit is `kbps` - `grace_period_millis` default set to 300 milliseconds. - `grace_period_timeout_secs` default set to 2 seconds. -- `geopip` is for resolving IP addresses to country names. +- `geoip` is for resolving IP addresses to country names. ##### 1.6.1.1 `retry` If set to `true` on connection loss to provider, the stream will be reconnected. diff --git a/backend/src/api/hdhomerun_proprietary.rs b/backend/src/api/hdhomerun_proprietary.rs index 1253d8f16..80d7178a0 100644 --- a/backend/src/api/hdhomerun_proprietary.rs +++ b/backend/src/api/hdhomerun_proprietary.rs @@ -3,7 +3,7 @@ use crate::model::{AppConfig, HdHomeRunDeviceConfig}; use bytes::{Buf, BufMut, BytesMut}; use log::{error, info, trace}; use std::collections::HashMap; -use std::io::{Cursor}; +use std::io::Cursor; use std::net::{Ipv4Addr, SocketAddr}; use std::sync::Arc; use std::time::Duration; @@ -36,13 +36,13 @@ mod packet { fn write_tlv_u32(buf: &mut BytesMut, tag: u8, value: u32) { buf.put_u8(tag); - buf.put_u8(4); + buf.put_u16(4); buf.put_u32(value); } fn write_tlv_u8(buf: &mut BytesMut, tag: u8, value: u8) { buf.put_u8(tag); - buf.put_u8(1); + buf.put_u16(1); buf.put_u8(value); } @@ -96,7 +96,11 @@ fn build_discover_response(device: &HdHomeRunDeviceConfig, server_host: &str) -> let mut response = BytesMut::new(); response.put_u16(packet::HDHOMERUN_TYPE_DISCOVER_RSP); - response.put_u16(u16::try_from(payload.len()).unwrap_or(0)); + //response.put_u16(u16::try_from(payload.len()).unwrap_or(0)); + if let Ok(n) = u16::try_from(payload.len()) { response.put_u16(n) } else { + error!("HDHR response payload too large ({} bytes)", payload.len()); + return Vec::new(); + } response.put(payload); let crc = crc32fast::hash(&response); @@ -109,7 +113,7 @@ fn parse_tlv(cursor: &mut Cursor<&[u8]>) -> HashMap> { let mut tags = HashMap::new(); loop { - let pos = match usize::try_from(cursor.position()) { + let pos = match usize::try_from(cursor.position()) { Ok(pos) => pos, Err(_err) => { return tags; @@ -170,9 +174,27 @@ async fn proprietary_discover_loop( continue; } + // let mut cursor = Cursor::new(data); + // let msg_type = cursor.get_u16(); + // let _msg_len = cursor.get_u16(); + let mut cursor = Cursor::new(data); let msg_type = cursor.get_u16(); - let _msg_len = cursor.get_u16(); + let msg_len = cursor.get_u16() as usize; + + // Validate total size: header (4) + payload (msg_len) + CRC (4) + if data.len() < 4 + msg_len + 4 { + trace!("Short HDHR discovery packet from {remote_addr}"); + continue; + } + let payload_end = 4 + msg_len; + let (framed, crc_tail) = data.split_at(payload_end); + let received_crc = u32::from_le_bytes(crc_tail[..4].try_into().unwrap_or_default()); + if crc32fast::hash(framed) != received_crc { + trace!("Invalid CRC in HDHR discovery packet from {remote_addr}"); + continue; + } + let mut cursor = Cursor::new(&framed[4..]); // parse TLV over payload only if msg_type == packet::HDHOMERUN_TYPE_DISCOVER_REQ { let tags = parse_tlv(&mut cursor); diff --git a/backend/src/auth/fingerprint.rs b/backend/src/auth/fingerprint.rs index de0b6a6af..cbe4e1a95 100644 --- a/backend/src/auth/fingerprint.rs +++ b/backend/src/auth/fingerprint.rs @@ -67,7 +67,7 @@ impl Fingerprint { let client_ip_port =format!("{client_ip}:{}", addr.port()); let ua = user_agent.unwrap_or_else(String::new); - let key = format!("{client_ip }|{ua}"); + let key = format!("{client_ip}|{ua}"); Ok(Fingerprint(key, client_ip_port)) } diff --git a/backend/src/processing/processor/playlist.rs b/backend/src/processing/processor/playlist.rs index 0f8ae8be4..9b55887b0 100644 --- a/backend/src/processing/processor/playlist.rs +++ b/backend/src/processing/processor/playlist.rs @@ -27,6 +27,7 @@ use crate::repository::playlist_repository::persist_playlist; use crate::utils::debug_if_enabled; use crate::utils::StepMeasure; use deunicode::deunicode; +use futures::StreamExt; use log::{debug, error, info, log_enabled, trace, warn, Level}; use reqwest::Client; use shared::error::{get_errors_notify_message, notify_err, TuliproxError}; @@ -627,15 +628,13 @@ async fn process_watch(cfg: &Config, client: &Arc, target: &Con return false; } - let mut futs = Vec::new(); - for pl in new_playlist { - if watches.iter().any(|r| r.is_match(&pl.title)) { - futs.push(process_group_watch(client, cfg, &target.name, pl)); - } - } - if !futs.is_empty() { - futures::future::join_all(futs).await; - } + futures::stream::iter( + new_playlist + .iter() + .filter(|pl| watches.iter().any(|r| r.is_match(&pl.title))) + .map(|pl| process_group_watch(client, cfg, &target.name, pl)) + ).for_each_concurrent(16, |f| f).await; + true } else { false diff --git a/backend/src/utils/geoip.rs b/backend/src/utils/geoip.rs index f5c6fa762..45605b955 100644 --- a/backend/src/utils/geoip.rs +++ b/backend/src/utils/geoip.rs @@ -34,8 +34,6 @@ impl GeoIp { if reader.read_line(&mut buf)? == 0 { break; } let line = buf.trim(); if line.is_empty() || line.starts_with('#') { continue; } - let line = buf.trim(); - if line.is_empty() || line.starts_with('#') { continue; } let parts: Vec<&str> = line.split(',').collect(); if parts.len() != 3 { continue; }