stream table ip fixed

This commit is contained in:
euzu
2025-10-30 13:00:25 +01:00
parent 5bf90fbd39
commit d71e8080f5
16 changed files with 452 additions and 125 deletions
+20 -19
View File
@@ -123,6 +123,7 @@ pub use try_option_bad_request;
pub use try_result_bad_request;
pub use try_result_not_found;
pub use try_unwrap_body;
use crate::auth::Fingerprint;
pub fn get_server_time() -> String {
chrono::offset::Local::now()
@@ -382,7 +383,7 @@ struct StreamingStrategy {
async fn resolve_streaming_strategy(
app_state: &AppState,
stream_url: &str,
addr: &str,
fingerprint: &Fingerprint,
input: &ConfigInput,
force_provider: Option<&str>,
) -> StreamingStrategy {
@@ -391,13 +392,13 @@ async fn resolve_streaming_strategy(
Some(provider) => {
app_state
.active_provider
.force_exact_acquire_connection(provider, addr)
.force_exact_acquire_connection(provider, &fingerprint.addr)
.await
}
None => {
app_state
.active_provider
.acquire_connection(&input.name, addr)
.acquire_connection(&input.name, &fingerprint.addr)
.await
}
};
@@ -461,7 +462,7 @@ async fn create_stream_response_details(
app_state: &AppState,
stream_options: &StreamOptions,
stream_url: &str,
addr: &str,
fingerprint: &Fingerprint,
req_headers: &HeaderMap,
input: &ConfigInput,
item_type: PlaylistItemType,
@@ -469,7 +470,7 @@ async fn create_stream_response_details(
connection_permission: UserConnectionPermission,
force_provider: Option<&str>,
) -> StreamDetails {
let mut streaming_strategy = resolve_streaming_strategy(app_state, stream_url, addr, input, force_provider).await;
let mut streaming_strategy = resolve_streaming_strategy(app_state, stream_url, fingerprint, input, force_provider).await;
let config_grace_period_millis = app_state
.app_config
.config
@@ -736,7 +737,7 @@ fn prepare_body_stream(
/// # Panics
pub async fn force_provider_stream_response(
addr: &str,
fingerprint: &Fingerprint,
app_state: &AppState,
user_session: &UserSession,
mut stream_channel: StreamChannel,
@@ -753,7 +754,7 @@ pub async fn force_provider_stream_response(
app_state,
&stream_options,
&user_session.stream_url,
addr,
fingerprint,
req_headers,
input,
item_type,
@@ -770,11 +771,11 @@ pub async fn force_provider_stream_response(
.map(|(h, sc, url)| (h.clone(), *sc, url.clone()));
app_state
.active_users
.update_session_addr(&user.username, &user_session.token, addr)
.update_session_addr(&user.username, &user_session.token, &fingerprint.addr)
.await;
stream_channel.shared = share_stream;
let stream =
ActiveClientStream::new(stream_details, app_state, user, connection_permission, addr, stream_channel, req_headers)
ActiveClientStream::new(stream_details, app_state, user, connection_permission, fingerprint, stream_channel, req_headers)
.await;
let (status_code, header_map) =
@@ -809,7 +810,7 @@ pub async fn force_provider_stream_response(
/// # Panics
#[allow(clippy::too_many_arguments, clippy::too_many_lines)]
pub async fn stream_response(
addr: &str,
fingerprint: &Fingerprint,
app_state: &AppState,
session_token: &str,
mut stream_channel: StreamChannel,
@@ -838,7 +839,7 @@ pub async fn stream_response(
let share_stream = is_stream_share_enabled(item_type, target);
if share_stream {
if let Some(value) =
shared_stream_response(app_state, stream_url, addr, user, connection_permission, stream_channel.clone(), req_headers).await
shared_stream_response(app_state, stream_url, fingerprint, user, connection_permission, stream_channel.clone(), req_headers).await
{
return value.into_response();
}
@@ -849,7 +850,7 @@ pub async fn stream_response(
app_state,
&stream_options,
stream_url,
addr,
fingerprint,
req_headers,
input,
item_type,
@@ -873,7 +874,7 @@ pub async fn stream_response(
};
stream_channel.shared = share_stream;
let stream =
ActiveClientStream::new(stream_details, app_state, user, connection_permission, addr, stream_channel, req_headers)
ActiveClientStream::new(stream_details, app_state, user, connection_permission, fingerprint, stream_channel, req_headers)
.await;
let stream_resp = if share_stream {
debug_if_enabled!(
@@ -889,7 +890,7 @@ pub async fn stream_response(
app_state,
stream_url,
stream,
Some(addr),
Some(&fingerprint.addr),
shared_headers,
stream_options.buffer_size,
provider_guard,
@@ -952,7 +953,7 @@ pub async fn stream_response(
virtual_id,
&provider,
&session_url,
addr,
&fingerprint.addr,
connection_permission,
)
.await;
@@ -984,14 +985,14 @@ fn get_stream_throttle(app_state: &AppState) -> u64 {
async fn shared_stream_response(
app_state: &AppState,
stream_url: &str,
addr: &str,
fingerprint: &Fingerprint,
user: &ProxyUserCredentials,
connect_permission: UserConnectionPermission,
mut stream_channel: StreamChannel,
req_headers: &HeaderMap,
) -> Option<impl IntoResponse> {
if let Some((stream, provider)) =
SharedStreamManager::subscribe_shared_stream(app_state, stream_url, Some(addr)).await
SharedStreamManager::subscribe_shared_stream(app_state, stream_url, Some(&fingerprint.addr)).await
{
debug_if_enabled!(
"Using shared stream {}",
@@ -1010,7 +1011,7 @@ async fn shared_stream_response(
stream_details.provider_name = provider;
stream_channel.shared = true;
let stream =
ActiveClientStream::new(stream_details, app_state, user, connect_permission, addr, stream_channel, req_headers)
ActiveClientStream::new(stream_details, app_state, user, connect_permission, fingerprint, stream_channel, req_headers)
.await
.boxed();
let mut response = axum::response::Response::builder().status(status_code);
@@ -1326,6 +1327,6 @@ pub fn json_or_bin_response<T: Serialize>(accept: Option<&String>, data: &T) ->
json_response(data).into_response()
}
pub fn create_fingerprint(fingerprint: &str, username: &str, virtual_id: u32) -> String {
pub fn create_session_fingerprint(fingerprint: &str, username: &str, virtual_id: u32) -> String {
format!("{fingerprint}|{username}|{virtual_id}")
}
+9 -11
View File
@@ -1,4 +1,4 @@
use crate::api::api_utils::{create_fingerprint, try_unwrap_body};
use crate::api::api_utils::{create_session_fingerprint, try_unwrap_body};
use crate::api::api_utils::{
force_provider_stream_response, get_stream_alternative_url, is_seek_request,
};
@@ -41,8 +41,7 @@ fn hls_response(hls_content: String) -> impl IntoResponse + Send {
#[allow(clippy::too_many_arguments)]
pub(in crate::api) async fn handle_hls_stream_request(
fingerprint: &str,
addr: &str,
fingerprint: &Fingerprint,
app_state: &Arc<AppState>,
user: &ProxyUserCredentials,
user_session: Option<&UserSession>,
@@ -59,7 +58,7 @@ pub(in crate::api) async fn handle_hls_stream_request(
Some(session) => {
match app_state
.active_provider
.force_exact_acquire_connection(&session.provider, addr)
.force_exact_acquire_connection(&session.provider, &fingerprint.addr)
.await
.get_provider_config()
{
@@ -78,14 +77,14 @@ pub(in crate::api) async fn handle_hls_stream_request(
{
Some(provider_cfg) => {
let stream_url = get_stream_alternative_url(&url, input, &provider_cfg);
let user_session_token = create_fingerprint(fingerprint, &user.username, virtual_id);
let user_session_token = create_session_fingerprint(&fingerprint.key, &user.username, virtual_id);
let session_token = app_state.active_users.create_user_session(
user,
&user_session_token,
virtual_id,
&provider_cfg.name,
&stream_url,
addr,
&fingerprint.addr,
connection_permission,
).await;
(stream_url, Some(session_token))
@@ -188,7 +187,7 @@ async fn resolve_stream_channel(
#[allow(clippy::too_many_lines)]
async fn hls_api_stream(
Fingerprint(fingerprint, addr): Fingerprint,
fingerprint: Fingerprint,
req_headers: axum::http::HeaderMap,
axum::extract::Path(params): axum::extract::Path<HlsApiPathParams>,
axum::extract::State(app_state): axum::extract::State<Arc<AppState>>,
@@ -218,7 +217,7 @@ async fn hls_api_stream(
)
);
let user_session_token = create_fingerprint(&fingerprint, &user.username, virtual_id);
let user_session_token = create_session_fingerprint(&fingerprint.key, &user.username, virtual_id);
let mut user_session = app_state
.active_users
.get_and_update_user_session(&user.username, &user_session_token).await;
@@ -258,7 +257,7 @@ async fn hls_api_stream(
if is_seek_request(stream_channel.cluster, &req_headers).await {
// partial request means we are in reverse proxy mode, seek happened
return force_provider_stream_response(
&addr,
&fingerprint,
&app_state,
session,
stream_channel,
@@ -285,7 +284,6 @@ async fn hls_api_stream(
if is_hls_url(&session.stream_url) {
return handle_hls_stream_request(
&fingerprint,
&addr,
&app_state,
&user,
Some(session),
@@ -301,7 +299,7 @@ async fn hls_api_stream(
let stream_channel = resolve_stream_channel(&app_state, &target, virtual_id, &hls_url).await;
force_provider_stream_response(
&addr,
&fingerprint,
&app_state,
session,
stream_channel,
+6 -9
View File
@@ -1,4 +1,4 @@
use crate::api::api_utils::{create_fingerprint, try_unwrap_body};
use crate::api::api_utils::{create_session_fingerprint, try_unwrap_body};
use crate::api::api_utils::{
force_provider_stream_response, get_user_target, get_user_target_by_credentials,
is_seek_request, redirect, redirect_response, resource_response, separate_number_and_remainder,
@@ -72,8 +72,7 @@ async fn m3u_api_post(
#[allow(clippy::too_many_lines)]
async fn m3u_api_stream(
fingerprint: &str,
addr: &str,
fingerprint: &Fingerprint,
req_headers: &axum::http::HeaderMap,
app_state: &Arc<AppState>,
api_req: &UserApiRequest,
@@ -124,7 +123,7 @@ async fn m3u_api_stream(
);
let cluster = XtreamCluster::try_from(pli.item_type).unwrap_or(XtreamCluster::Live);
let session_key = create_fingerprint(fingerprint, &user.username, virtual_id);
let session_key = create_session_fingerprint(&fingerprint.key, &user.username, virtual_id);
let user_session = app_state
.active_users
.get_and_update_user_session(&user.username, &session_key).await;
@@ -152,7 +151,7 @@ async fn m3u_api_stream(
if session.virtual_id == virtual_id && is_seek_request(cluster, req_headers).await {
// partial request means we are in reverse proxy mode, seek happened
return force_provider_stream_response(
addr,
fingerprint,
app_state,
session,
pli.to_stream_channel(),
@@ -208,7 +207,6 @@ async fn m3u_api_stream(
if is_hls_request {
return handle_hls_stream_request(
fingerprint,
addr,
app_state,
&user,
user_session.as_ref(),
@@ -223,7 +221,7 @@ async fn m3u_api_stream(
}
stream_response(
addr,
fingerprint,
app_state,
&session_key,
pli.to_stream_channel(),
@@ -302,7 +300,7 @@ async fn m3u_api_resource(
macro_rules! create_m3u_api_stream {
($fn_name:ident, $context:expr) => {
async fn $fn_name(
Fingerprint(fingerprint, addr): Fingerprint,
fingerprint: Fingerprint,
req_headers: axum::http::HeaderMap,
axum::extract::Query(api_req): axum::extract::Query<UserApiRequest>,
axum::extract::Path((username, password, stream_id)): axum::extract::Path<(
@@ -315,7 +313,6 @@ macro_rules! create_m3u_api_stream {
) -> impl IntoResponse + Send {
m3u_api_stream(
&fingerprint,
&addr,
&req_headers,
&app_state,
&api_req,
+3 -2
View File
@@ -8,7 +8,7 @@ use crate::auth::validator_admin;
use crate::utils::ip_checker::get_ips;
use crate::{VERSION};
use axum::response::IntoResponse;
use shared::model::{InputFetchMethod, IpCheckDto, StatusCheck};
use shared::model::{default_geoip_url, InputFetchMethod, IpCheckDto, StatusCheck};
use shared::utils::{concat_path_leading_slash};
use std::collections::{BTreeMap, HashMap};
use std::io::{Cursor};
@@ -80,8 +80,9 @@ async fn geoip_update(axum::extract::State(app_state): axum::extract::State<Arc<
let geoip_db_path = &*get_geoip_path(&config.working_dir);
let _file_lock = app_state.app_config.file_locks.write_lock(geoip_db_path);
let url = if geoip.url.trim().is_empty() { default_geoip_url() } else { geoip.url.clone() };
let input_source = InputSource {
url: geoip.url.clone(),
url,
username: None,
password: None,
method: InputFetchMethod::GET,
+20 -22
View File
@@ -1,7 +1,7 @@
// https://github.com/tellytv/go.xtream-codes/blob/master/structs.go
use crate::api::api_utils;
use crate::api::api_utils::{create_fingerprint, try_unwrap_body};
use crate::api::api_utils::{create_session_fingerprint, try_unwrap_body};
use crate::api::api_utils::{
force_provider_stream_response, get_user_target, get_user_target_by_credentials,
is_seek_request, redirect_response, resource_response, separate_number_and_remainder,
@@ -211,8 +211,7 @@ async fn get_user_info(user: &ProxyUserCredentials, app_state: &AppState) -> Xtr
#[allow(clippy::too_many_lines)]
async fn xtream_player_api_stream(
fingerprint: &str,
addr: &str,
fingerprint: &Fingerprint,
req_headers: &HeaderMap,
app_state: &Arc<AppState>,
api_req: &UserApiRequest,
@@ -269,7 +268,7 @@ async fn xtream_player_api_stream(
(pli.xtream_cluster, pli.item_type)
};
let session_key = create_fingerprint(fingerprint, &user.username, virtual_id);
let session_key = create_session_fingerprint(&fingerprint.key, &user.username, virtual_id);
let user_session = app_state
.active_users
.get_and_update_user_session(&user.username, &session_key).await;
@@ -295,13 +294,16 @@ async fn xtream_player_api_stream(
.into_response();
}
let mut stream_channel = pli.to_stream_channel();
stream_channel.item_type = item_type;
if session.virtual_id == virtual_id && is_seek_request(cluster, req_headers).await {
// partial request means we are in reverse proxy mode, seek happened
return force_provider_stream_response(
addr,
fingerprint,
app_state,
session,
pli.to_stream_channel(),
stream_channel,
req_headers,
&input,
&user,
@@ -369,7 +371,6 @@ async fn xtream_player_api_stream(
if is_hls_request {
return handle_hls_stream_request(
fingerprint,
addr,
app_state,
&user,
user_session.as_ref(),
@@ -383,11 +384,14 @@ async fn xtream_player_api_stream(
.into_response();
}
let mut stream_channel = pli.to_stream_channel();
stream_channel.item_type = item_type;
stream_response(
addr,
fingerprint,
app_state,
session_key.as_str(),
pli.to_stream_channel(),
stream_channel,
&stream_url,
req_headers,
&input,
@@ -402,8 +406,7 @@ async fn xtream_player_api_stream(
#[allow(clippy::too_many_lines)]
// Used by webui
async fn xtream_player_api_stream_with_token(
fingerprint: &str,
addr: &str,
fingerprint: &Fingerprint,
req_headers: &HeaderMap,
app_state: &Arc<AppState>,
target_id: u16,
@@ -439,7 +442,7 @@ async fn xtream_player_api_stream_with_token(
)
);
let session_key = create_fingerprint(fingerprint, "webui", virtual_id);
let session_key = create_session_fingerprint(&fingerprint.key, "webui", virtual_id);
let is_hls_request =
pli.item_type == PlaylistItemType::LiveHls || stream_ext.as_deref() == Some(HLS_EXT);
@@ -472,7 +475,6 @@ async fn xtream_player_api_stream_with_token(
if is_hls_request {
return handle_hls_stream_request(
fingerprint,
addr,
app_state,
&user,
None,
@@ -516,7 +518,7 @@ async fn xtream_player_api_stream_with_token(
sanitize_sensitive_info(&stream_url)
);
stream_response(
addr,
fingerprint,
app_state,
session_key.as_str(),
pli.to_stream_channel(),
@@ -770,7 +772,7 @@ async fn xtream_player_api_resource(
macro_rules! create_xtream_player_api_stream {
($fn_name:ident, $context:expr) => {
async fn $fn_name(
Fingerprint(fingerprint, addr): Fingerprint,
fingerprint: Fingerprint,
req_headers: HeaderMap,
axum::extract::Path((username, password, stream_id)): axum::extract::Path<(
String,
@@ -782,7 +784,6 @@ macro_rules! create_xtream_player_api_stream {
) -> impl IntoResponse + Send {
xtream_player_api_stream(
&fingerprint,
&addr,
&req_headers,
&app_state,
&api_req,
@@ -848,7 +849,7 @@ struct XtreamApiTimeShiftRequest {
}
async fn xtream_player_api_timeshift_stream(
Fingerprint(fingerprint, addr): Fingerprint,
fingerprint: Fingerprint,
req_headers: HeaderMap,
axum::extract::Query(mut api_req): axum::extract::Query<UserApiRequest>,
axum::extract::Path(timeshift_request): axum::extract::Path<XtreamApiTimeShiftRequest>,
@@ -891,7 +892,6 @@ async fn xtream_player_api_timeshift_stream(
xtream_player_api_stream(
&fingerprint,
&addr,
&req_headers,
&app_state,
&api_req,
@@ -908,7 +908,7 @@ async fn xtream_player_api_timeshift_stream(
}
async fn xtream_player_api_timeshift_query_stream(
Fingerprint(fingerprint, addr): Fingerprint,
fingerprint: Fingerprint,
req_headers: HeaderMap,
axum::extract::Query(api_query_req): axum::extract::Query<UserApiRequest>,
axum::extract::State(app_state): axum::extract::State<Arc<AppState>>,
@@ -933,7 +933,6 @@ async fn xtream_player_api_timeshift_query_stream(
}
xtream_player_api_stream(
&fingerprint,
&addr,
&req_headers,
&app_state,
&api_query_req,
@@ -1585,7 +1584,7 @@ macro_rules! register_xtream_api_timeshift {
}
async fn xtream_player_token_stream(
Fingerprint(fingerprint, addr): Fingerprint,
fingerprint: Fingerprint,
axum::extract::Path((token, target_id, cluster, stream_id)): axum::extract::Path<(
String,
u16,
@@ -1598,7 +1597,6 @@ async fn xtream_player_token_stream(
let ctxt = try_result_bad_request!(ApiStreamContext::from_str(cluster.as_str()));
xtream_player_api_stream_with_token(
&fingerprint,
&addr,
&req_headers,
&app_state,
target_id,
+7 -5
View File
@@ -14,6 +14,7 @@ use std::sync::Arc;
use arc_swap::ArcSwapOption;
use tokio::sync::mpsc::{unbounded_channel, UnboundedSender};
use tokio::sync::RwLock;
use crate::auth::Fingerprint;
use crate::utils::GeoIp;
const USER_GC_TTL: u64 = 900; // 15 Min
@@ -331,11 +332,11 @@ impl ActiveUserManager {
Self::get_active_connections(&self.user).await
}
pub async fn add_connection(&self, username: &str, max_connections: u32, addr: &str, provider: &str, stream_channel: StreamChannel, user_agent: Cow<'_, str>) -> UserConnectionGuard {
pub async fn add_connection(&self, username: &str, max_connections: u32, fingerprint: &Fingerprint, provider: &str, stream_channel: StreamChannel, user_agent: Cow<'_, str>) -> UserConnectionGuard {
let country = {
let geoip = self.geo_ip.load();
if let Some(geoip_db) = (*geoip).as_ref() {
geoip_db.lookup(&strip_port(addr))
geoip_db.lookup(&strip_port(&fingerprint.client_ip))
} else {
None
}
@@ -343,7 +344,8 @@ impl ActiveUserManager {
let stream_info = StreamInfo::new(
username,
addr,
&fingerprint.addr,
&fingerprint.client_ip,
provider,
stream_channel,
user_agent.to_string(),
@@ -364,7 +366,7 @@ impl ActiveUserManager {
{
let mut user_by_addr = self.user_by_addr.write().await;
user_by_addr.insert(addr.to_string(), username.to_string());
user_by_addr.insert(fingerprint.addr.to_string(), username.to_string());
}
if let Err(err) = self.connection_change_tx.send(ActiveUserConnectionChange::Connected(stream_info)) {
@@ -372,7 +374,7 @@ impl ActiveUserManager {
}
self.log_active_user().await;
UserConnectionGuard::new(Arc::new(self.clone_inner()), addr, self.release_sender())
UserConnectionGuard::new(Arc::new(self.clone_inner()), &fingerprint.addr, self.release_sender())
}
fn is_log_user_enabled(&self) -> bool {
@@ -18,6 +18,7 @@ use std::task::{Poll};
use axum::http::header::USER_AGENT;
use axum::http::HeaderMap;
use futures::task::AtomicWaker;
use crate::auth::Fingerprint;
const INNER_STREAM: u8 = 0_u8;
const GRACE_BLOCK_STREAM: u8 = 1_u8;
@@ -40,7 +41,7 @@ impl ActiveClientStream {
app_state: &AppState,
user: &ProxyUserCredentials,
connection_permission: UserConnectionPermission,
addr: &str,
fingerprint: &Fingerprint,
stream_channel: StreamChannel,
req_headers: &HeaderMap) -> Self {
if connection_permission == UserConnectionPermission::Exhausted {
@@ -51,11 +52,11 @@ impl ActiveClientStream {
let provider_name = stream_details.provider_name.as_ref().map_or_else(String::new, ToString::to_string);
let user_agent = req_headers.get(USER_AGENT).map(|h| String::from_utf8_lossy(h.as_bytes())).unwrap_or_default();
let user_connection_guard = Some(app_state.active_users.add_connection(username, user.max_connections, addr, &provider_name, stream_channel, user_agent).await);
let user_connection_guard = Some(app_state.active_users.add_connection(username, user.max_connections, fingerprint, &provider_name, stream_channel, user_agent).await);
let cfg = &app_state.app_config;
let waker = Some(Arc::new(AtomicWaker::new()));
let waker_clone = waker.clone();
let grace_stop_flag = Self::stream_grace_period(app_state, &stream_details, grant_user_grace_period, user, addr, waker_clone.clone());
let grace_stop_flag = Self::stream_grace_period(app_state, &stream_details, grant_user_grace_period, user, &fingerprint.addr, waker_clone.clone());
let custom_response = cfg.custom_stream_response.load();
let custom_video = custom_response.as_ref()
.map_or((None, None), |c|
+18 -2
View File
@@ -2,6 +2,7 @@ use std::net::SocketAddr;
use axum::extract::{ConnectInfo, FromRequestParts};
use axum::http::request::Parts;
use axum::http::StatusCode;
use log::debug;
use crate::auth::Rejection;
const MAX_HEADER_LENGTH: usize = 512;
@@ -16,8 +17,21 @@ fn validate_header(value: &str) -> Option<String> {
}
#[derive(Debug, PartialEq, Eq, Clone)]
pub struct Fingerprint(pub String, pub String);
pub struct Fingerprint {
pub key: String,
pub client_ip: String,
pub addr: String,
}
impl Fingerprint {
pub fn new(key: String, client_ip: String, addr: String) -> Self {
Self {
key,
client_ip,
addr,
}
}
}
impl<B> FromRequestParts<B> for Fingerprint
where
@@ -67,6 +81,8 @@ impl Fingerprint {
let ua = user_agent.unwrap_or_else(String::new);
let key = format!("{client_ip}|{ua}");
Ok(Fingerprint(key, addr.to_string()))
debug!("{key}, {client_ip}, {addr}");
Ok(Fingerprint::new(key, client_ip, addr.to_string()))
}
}
+30 -14
View File
@@ -1,9 +1,9 @@
use crate::repository::bplustree::BPlusTree;
use serde::{Deserialize, Serialize};
use std::io;
use std::io::{BufRead};
use std::io::BufRead;
use std::net::Ipv4Addr;
use std::path::Path;
use serde::{Serialize, Deserialize};
use crate::repository::bplustree::BPlusTree;
fn ipv4_to_u32(ip: &str) -> Option<u32> {
ip.parse::<Ipv4Addr>().ok().map(u32::from)
@@ -16,14 +16,31 @@ pub struct GeoIp {
impl GeoIp {
pub fn load(path: &Path) -> io::Result<Self> {
let tree = BPlusTree::load(path)?;
Ok(Self { tree })
let mut tree = BPlusTree::load(path)?;
let private_ranges = vec![
("127.0.0.0", "127.255.255.255", "Loopback"),
("10.0.0.0", "10.255.255.255", "LAN"),
("172.16.0.0", "172.31.255.255", "LAN"),
("192.168.0.0", "192.168.255.255", "LAN"),
("169.254.0.0", "169.254.255.255", "Link-Local"),
("172.17.0.0", "172.17.255.255", "Docker"),
("172.18.0.0", "172.31.255.255", "Docker")
];
for range in private_ranges {
if let (Some(start), Some(end)) = (ipv4_to_u32(range.0), ipv4_to_u32(range.1)) {
let cc = range.2.to_string();
tree.insert(start, (end, cc));
}
}
Ok(Self { tree })
}
pub fn new() -> Self {
Self { tree: BPlusTree::new() }
Self { tree: BPlusTree::new() }
}
pub fn import_ipv4_from_csv(&mut self, mut reader: impl BufRead, db_path: &Path) -> std::io::Result<u64> {
@@ -59,23 +76,23 @@ impl GeoIp {
}
impl Default for GeoIp {
fn default() -> Self {
Self::new()
}
fn default() -> Self {
Self::new()
}
}
#[cfg(test)]
mod test {
// https://raw.githubusercontent.com/sapics/ip-location-db/refs/heads/main/asn-country/asn-country-ipv4.csv
use crate::utils::geoip::GeoIp;
use std::fs::File;
use std::io::BufReader;
use std::path::PathBuf;
use crate::utils::geoip::GeoIp;
#[test]
pub fn test_csv() {
let db_file = PathBuf::from("/projects/m3u-test/asn-country-ipv4.db");
let db_file = PathBuf::from("/projects/m3u-test/asn-country-ipv4.db");
let source = PathBuf::from("/projects/m3u-test/asn-country-ipv4.csv");
let file = File::open(source).expect("Could not open csv file");
let reader = BufReader::new(file);
@@ -83,11 +100,10 @@ mod test {
let _ = geo_ip.import_ipv4_from_csv(reader, &db_file).expect("Could not import csv");
let geo_ip = GeoIp::load(&db_file).expect("Failed to load geoip db");
if let Some(cc) = geo_ip.lookup("72.13.24.23") {
if let Some(cc) = geo_ip.lookup("72.13.24.23") {
assert_eq!(cc, "US");
} else {
assert!(false);
}
}
}
+255
View File
@@ -405,5 +405,260 @@
"MSG": {
"WELCOME": "Welcome to tuliprox setup. This wizard will guide you through the initial setup process."
}
},
"COUNTRY" : {
"Loopback": "Loopback",
"LAN":"LAN",
"Link-Local":"Link-Local",
"Docker":"Docker",
"AF": "Afghanistan",
"AX": "Åland Islands",
"AL": "Albania",
"DZ": "Algeria",
"AS": "American Samoa",
"AD": "Andorra",
"AO": "Angola",
"AI": "Anguilla",
"AQ": "Antarctica",
"AG": "Antigua and Barbuda",
"AR": "Argentina",
"AM": "Armenia",
"AW": "Aruba",
"AU": "Australia",
"AT": "Austria",
"AZ": "Azerbaijan",
"BS": "Bahamas",
"BH": "Bahrain",
"BD": "Bangladesh",
"BB": "Barbados",
"BY": "Belarus",
"BE": "Belgium",
"BZ": "Belize",
"BJ": "Benin",
"BM": "Bermuda",
"BT": "Bhutan",
"BO": "Bolivia",
"BQ": "Bonaire, Sint Eustatius and Saba",
"BA": "Bosnia and Herzegovina",
"BW": "Botswana",
"BV": "Bouvet Island",
"BR": "Brazil",
"IO": "British Indian Ocean Territory",
"BN": "Brunei Darussalam",
"BG": "Bulgaria",
"BF": "Burkina Faso",
"BI": "Burundi",
"KH": "Cambodia",
"CM": "Cameroon",
"CA": "Canada",
"CV": "Cabo Verde",
"KY": "Cayman Islands",
"CF": "Central African Republic",
"TD": "Chad",
"CL": "Chile",
"CN": "China",
"CX": "Christmas Island",
"CC": "Cocos (Keeling) Islands",
"CO": "Colombia",
"KM": "Comoros",
"CG": "Congo",
"CD": "Congo, Democratic Republic of the",
"CK": "Cook Islands",
"CR": "Costa Rica",
"CI": "Côte d'Ivoire",
"HR": "Croatia",
"CU": "Cuba",
"CW": "Curaçao",
"CY": "Cyprus",
"CZ": "Czechia",
"DK": "Denmark",
"DJ": "Djibouti",
"DM": "Dominica",
"DO": "Dominican Republic",
"EC": "Ecuador",
"EG": "Egypt",
"SV": "El Salvador",
"GQ": "Equatorial Guinea",
"ER": "Eritrea",
"EE": "Estonia",
"SZ": "Eswatini",
"ET": "Ethiopia",
"FK": "Falkland Islands (Malvinas)",
"FO": "Faroe Islands",
"FJ": "Fiji",
"FI": "Finland",
"FR": "France",
"GF": "French Guiana",
"PF": "French Polynesia",
"TF": "French Southern Territories",
"GA": "Gabon",
"GM": "Gambia",
"GE": "Georgia",
"DE": "Germany",
"GH": "Ghana",
"GI": "Gibraltar",
"GR": "Greece",
"GL": "Greenland",
"GD": "Grenada",
"GP": "Guadeloupe",
"GU": "Guam",
"GT": "Guatemala",
"GG": "Guernsey",
"GN": "Guinea",
"GW": "Guinea-Bissau",
"GY": "Guyana",
"HT": "Haiti",
"HM": "Heard Island and McDonald Islands",
"VA": "Holy See",
"HN": "Honduras",
"HK": "Hong Kong",
"HU": "Hungary",
"IS": "Iceland",
"IN": "India",
"ID": "Indonesia",
"IR": "Iran",
"IQ": "Iraq",
"IE": "Ireland",
"IM": "Isle of Man",
"IL": "Israel",
"IT": "Italy",
"JM": "Jamaica",
"JP": "Japan",
"JE": "Jersey",
"JO": "Jordan",
"KZ": "Kazakhstan",
"KE": "Kenya",
"KI": "Kiribati",
"KP": "Korea (North)",
"KR": "Korea (South)",
"KW": "Kuwait",
"KG": "Kyrgyzstan",
"LA": "Lao People's Democratic Republic",
"LV": "Latvia",
"LB": "Lebanon",
"LS": "Lesotho",
"LR": "Liberia",
"LY": "Libya",
"LI": "Liechtenstein",
"LT": "Lithuania",
"LU": "Luxembourg",
"MO": "Macao",
"MG": "Madagascar",
"MW": "Malawi",
"MY": "Malaysia",
"MV": "Maldives",
"ML": "Mali",
"MT": "Malta",
"MH": "Marshall Islands",
"MQ": "Martinique",
"MR": "Mauritania",
"MU": "Mauritius",
"YT": "Mayotte",
"MX": "Mexico",
"FM": "Micronesia",
"MD": "Moldova",
"MC": "Monaco",
"MN": "Mongolia",
"ME": "Montenegro",
"MS": "Montserrat",
"MA": "Morocco",
"MZ": "Mozambique",
"MM": "Myanmar",
"NA": "Namibia",
"NR": "Nauru",
"NP": "Nepal",
"NL": "Netherlands",
"NC": "New Caledonia",
"NZ": "New Zealand",
"NI": "Nicaragua",
"NE": "Niger",
"NG": "Nigeria",
"NU": "Niue",
"NF": "Norfolk Island",
"MK": "North Macedonia",
"MP": "Northern Mariana Islands",
"NO": "Norway",
"OM": "Oman",
"PK": "Pakistan",
"PW": "Palau",
"PS": "Palestine, State of",
"PA": "Panama",
"PG": "Papua New Guinea",
"PY": "Paraguay",
"PE": "Peru",
"PH": "Philippines",
"PN": "Pitcairn",
"PL": "Poland",
"PT": "Portugal",
"PR": "Puerto Rico",
"QA": "Qatar",
"RE": "Réunion",
"RO": "Romania",
"RU": "Russian Federation",
"RW": "Rwanda",
"BL": "Saint Barthélemy",
"SH": "Saint Helena, Ascension and Tristan da Cunha",
"KN": "Saint Kitts and Nevis",
"LC": "Saint Lucia",
"MF": "Saint Martin (French part)",
"PM": "Saint Pierre and Miquelon",
"VC": "Saint Vincent and the Grenadines",
"WS": "Samoa",
"SM": "San Marino",
"ST": "Sao Tome and Principe",
"SA": "Saudi Arabia",
"SN": "Senegal",
"RS": "Serbia",
"SC": "Seychelles",
"SL": "Sierra Leone",
"SG": "Singapore",
"SX": "Sint Maarten (Dutch part)",
"SK": "Slovakia",
"SI": "Slovenia",
"SB": "Solomon Islands",
"SO": "Somalia",
"ZA": "South Africa",
"GS": "South Georgia and the South Sandwich Islands",
"SS": "South Sudan",
"ES": "Spain",
"LK": "Sri Lanka",
"SD": "Sudan",
"SR": "Suriname",
"SJ": "Svalbard and Jan Mayen",
"SE": "Sweden",
"CH": "Switzerland",
"SY": "Syrian Arab Republic",
"TW": "Taiwan",
"TJ": "Tajikistan",
"TZ": "Tanzania",
"TH": "Thailand",
"TL": "Timor-Leste",
"TG": "Togo",
"TK": "Tokelau",
"TO": "Tonga",
"TT": "Trinidad and Tobago",
"TN": "Tunisia",
"TR": "Türkiye",
"TM": "Turkmenistan",
"TC": "Turks and Caicos Islands",
"TV": "Tuvalu",
"UG": "Uganda",
"UA": "Ukraine",
"AE": "United Arab Emirates",
"GB": "United Kingdom",
"US": "United States",
"UM": "United States Minor Outlying Islands",
"UY": "Uruguay",
"UZ": "Uzbekistan",
"VU": "Vanuatu",
"VE": "Venezuela",
"VN": "Viet Nam",
"VG": "Virgin Islands (British)",
"VI": "Virgin Islands (U.S.)",
"WF": "Wallis and Futuna",
"EH": "Western Sahara",
"YE": "Yemen",
"ZM": "Zambia",
"ZW": "Zimbabwe"
}
}
@@ -1,33 +1,35 @@
use crate::app::components::menu_item::MenuItem;
use crate::app::components::popup_menu::PopupMenu;
use crate::app::components::{AppIcon, Table, TableDefinition, ToggleSwitch};
use crate::app::ConfigContext;
use crate::hooks::use_service_context;
use gloo_timers::callback::Interval;
use gloo_utils::window;
use shared::error::{create_tuliprox_error_result, TuliproxError, TuliproxErrorKind};
use shared::model::{SortOrder, StreamInfo};
use shared::utils::{current_time_secs, strip_port};
use std::fmt::Display;
use std::rc::Rc;
use std::str::FromStr;
use gloo_timers::callback::Interval;
use gloo_utils::window;
use wasm_bindgen::JsCast;
use web_sys::Element;
use yew::prelude::*;
use yew_i18n::use_translation;
use shared::utils::{current_time_secs, strip_port};
use crate::utils::t_safe;
const HEADERS: [&str; 12] = [
"LABEL.EMPTY",
"LABEL.USERNAME",
"LABEL.STREAM_ID",
"LABEL.CLUSTER",
"LABEL.CHANNEL",
"LABEL.GROUP",
"LABEL.CLIENT_IP",
"LABEL.COUNTRY",
"LABEL.PROVIDER",
"LABEL.SHARED",
"LABEL.USER_AGENT",
"LABEL.DURATION"
"EMPTY",
"USERNAME",
"STREAM_ID",
"CLUSTER",
"CHANNEL",
"GROUP",
"CLIENT_IP",
"COUNTRY",
"PROVIDER",
"SHARED",
"USER_AGENT",
"DURATION"
];
fn format_duration(seconds: u64) -> String {
@@ -62,10 +64,29 @@ pub struct StreamsTableProps {
pub fn StreamsTable(props: &StreamsTableProps) -> Html {
let translate = use_translation();
let services = use_service_context();
let config_ctx = use_context::<ConfigContext>().expect("Config context not found");
let popup_anchor_ref = use_state(|| None::<web_sys::Element>);
let popup_is_open = use_state(|| false);
let selected_dto = use_state(|| None::<Rc<StreamInfo>>);
let headers = use_memo(config_ctx, |cfg| {
let include_country = if let Some(app_cfg) = &cfg.config {
app_cfg.config.is_geoip_enabled()
} else {
false
};
let visible_headers: Vec<&str> = if include_country {
HEADERS.to_vec() // alle Header
} else {
HEADERS.iter()
.filter(|h| **h != "COUNTRY")
.copied()
.collect()
};
visible_headers
});
use_effect_with((), move |_| {
let interval = Interval::new(1000, update_timestamps);
@@ -95,11 +116,12 @@ pub fn StreamsTable(props: &StreamsTableProps) -> Html {
let render_header_cell = {
let translator = translate.clone();
let headers = headers.clone();
Callback::<usize, Html>::from(move |col| {
html! {
{
if col < HEADERS.len() {
translator.t(HEADERS[col])
if col < headers.len() {
translator.t(&format!("LABEL.{}", headers[col]))
} else {
String::new()
}
@@ -110,10 +132,12 @@ pub fn StreamsTable(props: &StreamsTableProps) -> Html {
let render_data_cell = {
let popup_onclick = handle_popup_onclick.clone();
let headers = headers.clone();
let translate = translate.clone();
Callback::<(usize, usize, Rc<StreamInfo>), Html>::from(
move |(row, col, dto): (usize, usize, Rc<StreamInfo>)| {
match col {
0 => {
match headers[col] {
"EMPTY" => {
let popup_onclick = popup_onclick.clone();
html! {
<button class="tp__icon-button"
@@ -123,22 +147,22 @@ pub fn StreamsTable(props: &StreamsTableProps) -> Html {
</button>
}
}
1 => html! {dto.username.as_str()},
2 => html! { <>
"USERNAME" => html! {dto.username.as_str()},
"STREAM_ID" => html! { <>
{ dto.channel.virtual_id.to_string() }
{" ("}
{ dto.channel.provider_id.to_string() }
{")"}
</>},
3 => html! {dto.channel.cluster},
4 => html! {dto.channel.title.as_str()},
5 => html! {dto.channel.group.as_str()},
6 => html! { strip_port(&dto.addr)},
7 => html! { dto.country.as_ref().map_or_else(String::new, |c| c.clone()) },
8 => html! {dto.provider.as_str()},
9 => html! { <ToggleSwitch value={dto.channel.shared} readonly={true} /> },
10 => html! { dto.user_agent.as_str() },
11 => html! { <span class="tp__stream-table__duration" data-ts={dto.ts.to_string()}>{format_duration(dto.ts)}</span> },
"CLUSTER" => html! {dto.channel.cluster},
"CHANNEL" => html! {dto.channel.title.as_str()},
"GROUP" => html! {dto.channel.group.as_str()},
"CLIENT_IP" => html! { strip_port(&dto.client_ip)},
"COUNTRY" => html! { dto.country.as_ref().map_or_else(String::new, |c| t_safe(&translate, &format!("COUNTRY.{c}"))) },
"PROVIDER" => html! {dto.provider.as_str()},
"SHARED" => html! { <ToggleSwitch value={dto.channel.shared} readonly={true} /> },
"USER_AGENT" => html! { dto.user_agent.as_str() },
"DURATION" => html! { <span class="tp__stream-table__duration" data-ts={dto.ts.to_string()}>{format_duration(dto.ts)}</span> },
_ => html! {""},
}
})
@@ -148,8 +172,7 @@ pub fn StreamsTable(props: &StreamsTableProps) -> Html {
false
});
let on_sort = Callback::<Option<(usize, SortOrder)>, ()>::from(move |_args| {
});
let on_sort = Callback::<Option<(usize, SortOrder)>, ()>::from(move |_args| {});
let table_definition = {
// first register for config update
@@ -157,11 +180,11 @@ pub fn StreamsTable(props: &StreamsTableProps) -> Html {
let render_data_cell_cb = render_data_cell.clone();
let is_sortable = is_sortable.clone();
let on_sort = on_sort.clone();
let num_cols = HEADERS.len();
let num_cols = headers.len();
use_memo(props.streams.clone(), move |streams| {
streams.as_ref().map(|list|
Rc::new(TableDefinition::<StreamInfo> {
items: if list.is_empty() {None} else {Some(Rc::new(list.clone()))},
items: if list.is_empty() { None } else { Some(Rc::new(list.clone())) },
num_cols,
is_sortable,
on_sort,
@@ -1,4 +1,5 @@
use gloo_timers::callback::Interval;
use gloo_utils::window;
use yew::prelude::*;
#[derive(Clone, Debug, Properties, PartialEq)]
@@ -14,8 +15,12 @@ struct Star {
#[function_component]
pub fn FloatingBackground() -> Html {
let stars = use_state(|| {
let width = window().inner_width().unwrap().as_f64().unwrap_or(800.0);
let height = window().inner_height().unwrap().as_f64().unwrap_or(600.0);
let area = width * height;
let num_stars = ((area / 50000.0) as usize).clamp(10, 40);
let mut rng = fastrand::Rng::new();
let stars: Vec<Star> = (0..40).map(|_| {
let stars: Vec<Star> = (0..num_stars).map(|_| {
Star {
x: rng.f64() * 100.0,
y: rng.f64() * 100.0,
@@ -31,7 +36,7 @@ pub fn FloatingBackground() -> Html {
{
let stars = stars.clone();
use_effect(move || {
let interval = Interval::new(50, move || {
let interval = Interval::new(200, move || {
stars.set(
stars.iter().map(|s| {
let mut star = s.clone();
+11
View File
@@ -3,6 +3,7 @@ mod storage;
use wasm_bindgen::JsCast;
use wasm_bindgen::prelude::Closure;
use web_sys::window;
use yew_i18n::YewI18n;
pub use storage::*;
#[macro_export]
@@ -30,4 +31,14 @@ where
millis,
)
.unwrap();
}
pub fn t_safe(i18n: &YewI18n, key: &str) -> String {
let result = i18n.t(key);
if result.starts_with("Unable to find the key") {
key.to_string()
} else {
result
}
}
@@ -1,5 +1,6 @@
use crate::model::StreamInfo;
#[allow(clippy::large_enum_variant)]
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, PartialEq)]
#[serde(tag = "type", content = "payload", rename_all = "camelCase")]
pub enum ActiveUserConnectionChange {
+1 -1
View File
@@ -1,5 +1,5 @@
fn default_geoip_url() -> String { String::from("https://raw.githubusercontent.com/sapics/ip-location-db/refs/heads/main/asn-country/asn-country-ipv4.csv") }
pub fn default_geoip_url() -> String { String::from("https://raw.githubusercontent.com/sapics/ip-location-db/refs/heads/main/asn-country/asn-country-ipv4.csv") }
#[derive(Debug, Clone, serde::Serialize, serde::Deserialize, Default, PartialEq)]
#[serde(deny_unknown_fields)]
+3 -1
View File
@@ -49,6 +49,7 @@ pub struct StreamInfo {
pub channel: StreamChannel,
pub provider: String,
pub addr: String,
pub client_ip: String,
#[serde(default)]
pub user_agent: String,
#[serde(default)]
@@ -58,12 +59,13 @@ pub struct StreamInfo {
}
impl StreamInfo {
pub fn new(username: &str, addr: &str, provider: &str, stream_channel: StreamChannel, user_agent: String, country: Option<String>) -> Self {
pub fn new(username: &str, addr: &str, client_ip: &str, provider: &str, stream_channel: StreamChannel, user_agent: String, country: Option<String>) -> Self {
Self {
username: username.to_string(),
channel: stream_channel,
provider: provider.to_string(),
addr: addr.to_string(),
client_ip: client_ip.to_string(),
user_agent,
ts: current_time_secs(),
country,