mirror of
https://github.com/euzu/tuliprox.git
synced 2026-10-04 15:02:16 +02:00
Some refactorings
- HdHomerun - WebUI GeoIP - Playlist Updates ...
This commit is contained in:
@@ -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.
|
||||
|
||||
@@ -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<u8, Vec<u8>> {
|
||||
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);
|
||||
|
||||
@@ -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))
|
||||
}
|
||||
|
||||
@@ -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<reqwest::Client>, 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
|
||||
|
||||
@@ -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; }
|
||||
|
||||
Reference in New Issue
Block a user