refactored user bouquet filtering

This commit is contained in:
euzu
2025-03-08 12:27:06 +01:00
parent 5721587019
commit c6b4463035
6 changed files with 70 additions and 59 deletions
+2 -2
View File
@@ -228,8 +228,8 @@ async fn lineup_json(app_state: web::Data<HdHomerunAppState>) -> impl Responder
None
};
let live_channels = XtreamPlaylistIterator::new(XtreamCluster::Live, &cfg, target, 0, &credentials).await.ok();
let vod_channels = XtreamPlaylistIterator::new(XtreamCluster::Video, &cfg, target, 0, &credentials).await.ok();
let live_channels = XtreamPlaylistIterator::new(XtreamCluster::Live, &cfg, target, None, &credentials).await.ok();
let vod_channels = XtreamPlaylistIterator::new(XtreamCluster::Video, &cfg, target, None, &credentials).await.ok();
// TODO include when resolved
//let series_channels = xtream_repository::iter_raw_xtream_playlist(cfg, target, XtreamCluster::Series);
let user_credentials = Arc::new(credentials);
+3 -3
View File
@@ -504,7 +504,7 @@ async fn xtream_get_short_epg(app_state: &AppState, user: &ProxyUserCredentials,
get_empty_epg_response()
}
async fn xtream_player_api_handle_content_action(config: &Config, target_name: &str, action: &str, category_id: &str, user: &ProxyUserCredentials, req: &HttpRequest) -> Option<HttpResponse> {
async fn xtream_player_api_handle_content_action(config: &Config, target_name: &str, action: &str, category_id: Option<u32>, user: &ProxyUserCredentials, req: &HttpRequest) -> Option<HttpResponse> {
if let Ok((path, content)) = match action {
ACTION_GET_LIVE_CATEGORIES => xtream_repository::xtream_get_collection_path(config, target_name, xtream_repository::COL_CAT_LIVE),
ACTION_GET_VOD_CATEGORIES => xtream_repository::xtream_get_collection_path(config, target_name, xtream_repository::COL_CAT_VOD),
@@ -624,14 +624,14 @@ async fn xtream_player_api(
_ => {}
}
let category_id = api_req.category_id.trim().parse::<u32>().ok();
// Handle general content actions
if let Some(response) = xtream_player_api_handle_content_action(
&app_state.config, &target.name, action, api_req.category_id.trim(), &user, req,
&app_state.config, &target.name, action, category_id, &user, req,
).await {
return response;
}
let category_id = api_req.category_id.trim().parse::<u32>().unwrap_or(0);
let result = match action {
ACTION_GET_LIVE_STREAMS =>
skip_flag_optional!(skip_live, xtream_repository::xtream_load_rewrite_playlist(XtreamCluster::Live, &app_state.config, target, category_id, &user).await),
+1 -1
View File
@@ -46,7 +46,7 @@ impl M3uPlaylistIterator {
let target_options = target.options.as_ref();
let include_type_in_url = target_options.is_some_and(|opts| opts.m3u_include_type_in_url);
let mask_redirect_url = target_options.is_some_and(|opts| opts.m3u_mask_redirect_url);
let filter = user_get_bouquet_filter(cfg, &user.username, "", TargetType::M3u, XtreamCluster::Live).await;
let filter = user_get_bouquet_filter(cfg, &user.username, None, TargetType::M3u, XtreamCluster::Live).await;
// TODO m3u bouquet filter
let server_info = cfg.get_user_server_info(user);
+60 -48
View File
@@ -85,7 +85,7 @@ fn add_target_user_to_user_tree(target_users: &[TargetUser], user_tree: &mut BPl
}
}
pub fn merge_api_user(cfg: &Config, target_users: &[TargetUser]) -> Result<u64, std::io::Error> {
pub fn merge_api_user(cfg: &Config, target_users: &[TargetUser]) -> Result<u64, Error> {
let path = get_api_user_db_path(cfg);
let lock = cfg.file_locks.read_lock(&path);
let mut user_tree: BPlusTree<String, StoredProxyUserCredentials> = BPlusTree::load(&path).unwrap_or_else(|_| BPlusTree::new());
@@ -108,7 +108,7 @@ pub fn backup_api_user_db_file(cfg: &Config, path: &Path) {
}
}
pub fn store_api_user(cfg: &Config, target_users: &[TargetUser]) -> Result<u64, std::io::Error> {
pub fn store_api_user(cfg: &Config, target_users: &[TargetUser]) -> Result<u64, Error> {
let mut user_tree = BPlusTree::<String, StoredProxyUserCredentials>::new();
add_target_user_to_user_tree(target_users, &mut user_tree);
let path = get_api_user_db_path(cfg);
@@ -117,7 +117,7 @@ pub fn store_api_user(cfg: &Config, target_users: &[TargetUser]) -> Result<u64,
user_tree.store(&path)
}
pub fn load_api_user(cfg: &Config) -> Result<Vec<TargetUser>, std::io::Error> {
pub fn load_api_user(cfg: &Config) -> Result<Vec<TargetUser>, Error> {
let path = get_api_user_db_path(cfg);
let lock = cfg.file_locks.read_lock(&path);
let user_tree = BPlusTree::<String, StoredProxyUserCredentials>::load(&path)?;
@@ -175,17 +175,33 @@ async fn save_xtream_user_bouquet_for_target(config: &Config, target_name: &str,
XtreamCluster::Video => user_get_vod_bouquet_path(storage_path, &TargetType::Xtream),
XtreamCluster::Series => user_get_series_bouquet_path(storage_path, &TargetType::Xtream),
};
match bouquet {
Some(bouquet) => {
if let Some(xtream_categories) = xtream_get_playlist_categories(config, target_name, cluster).await {
let filtered: Vec<&PlaylistXtreamCategory> = xtream_categories.iter().filter(|p| bouquet.contains(&p.name)).collect();
json_write_documents_to_file(&bouquet_path, &filtered)?;
} else {
json_write_documents_to_file(&bouquet_path, &None::<String>)?;
}
if let Some(bouquet_categories) = bouquet {
if let Some(xtream_categories) = xtream_get_playlist_categories(config, target_name, cluster).await {
let filtered: Vec<&PlaylistXtreamCategory> = xtream_categories.iter().filter(|p| bouquet_categories.contains(&p.name)).collect();
return json_write_documents_to_file(&bouquet_path, &filtered);
}
None => {
json_write_documents_to_file(&bouquet_path, &bouquet)?;
}
if bouquet_path.exists() {
std::fs::remove_file(bouquet_path)?;
}
Ok(())
}
fn save_m3u_user_bouquet_for_target(storage_path: &Path, target: &TargetType, cluster: XtreamCluster, bouquet: Option<&Vec<String>>) -> Result<(), Error> {
let bouquet_path = match cluster {
XtreamCluster::Live => user_get_live_bouquet_path(storage_path, target),
XtreamCluster::Video => user_get_vod_bouquet_path(storage_path, target),
XtreamCluster::Series => user_get_series_bouquet_path(storage_path, target),
};
match bouquet {
Some(bouquet_categories) => {
json_write_documents_to_file(&bouquet_path, bouquet_categories)?;
}
None => if bouquet_path.exists() {
std::fs::remove_file(bouquet_path)?;
}
}
@@ -198,14 +214,14 @@ async fn save_user_bouquet_for_target(config: &Config, target_name: &str, storag
save_xtream_user_bouquet_for_target(config, target_name, storage_path, XtreamCluster::Video, bouquet.vod.as_ref()).await?;
save_xtream_user_bouquet_for_target(config, target_name, storage_path, XtreamCluster::Series, bouquet.series.as_ref()).await?;
} else {
json_write_documents_to_file(&user_get_live_bouquet_path(storage_path, &target), &bouquet.live)?;
json_write_documents_to_file(&user_get_vod_bouquet_path(storage_path, &target), &bouquet.vod)?;
json_write_documents_to_file(&user_get_series_bouquet_path(storage_path, &target), &bouquet.series)?;
save_m3u_user_bouquet_for_target(storage_path, &target, XtreamCluster::Live, bouquet.live.as_ref())?;
save_m3u_user_bouquet_for_target(storage_path, &target, XtreamCluster::Video, bouquet.vod.as_ref())?;
save_m3u_user_bouquet_for_target(storage_path, &target, XtreamCluster::Series, bouquet.series.as_ref())?;
}
Ok(())
}
pub async fn save_user_bouquet(cfg: &Config, target_name: &str, username: &str, bouquet: &PlaylistBouquetDto) -> Result<(), std::io::Error> {
pub async fn save_user_bouquet(cfg: &Config, target_name: &str, username: &str, bouquet: &PlaylistBouquetDto) -> Result<(), Error> {
if let Some(storage_path) = ensure_user_storage_path(cfg, username) {
if let Some(xb) = &bouquet.xtream {
save_user_bouquet_for_target(cfg, target_name, &storage_path, TargetType::Xtream, xb).await?;
@@ -215,12 +231,12 @@ pub async fn save_user_bouquet(cfg: &Config, target_name: &str, username: &str,
}
Ok(())
} else {
Err(std::io::Error::new(std::io::ErrorKind::NotFound, format!("User config path not found for user {username}")))
Err(Error::new(std::io::ErrorKind::NotFound, format!("User config path not found for user {username}")))
}
}
async fn load_user_bouquet_from_file(file: &Path) -> Option<String> {
tokio::fs::read_to_string(file).await.ok()
tokio::fs::read_to_string(file).await.ok().filter(|content| !(content.is_empty() || content == "null"))
}
fn convert_xtream_user_bouquet(bouquet_cluster: Option<String>) -> Option<String> {
@@ -278,42 +294,38 @@ pub(crate) async fn user_get_series_bouquet(cfg: &Config, username: &str, target
user_get_cluster_bouquet(cfg, username, target, XtreamCluster::Series).await
}
pub async fn user_get_bouquet_filter(config: &Config, username: &str, category_id: &str, target: TargetType, cluster: XtreamCluster) -> Option<HashSet<String>> {
pub async fn user_get_bouquet_filter(config: &Config, username: &str, category_id: Option<u32>, target: TargetType, cluster: XtreamCluster) -> Option<HashSet<String>> {
if let Some(cid) = category_id {
return Some(HashSet::from([cid.to_string()]));
}
let bouquet = match cluster {
XtreamCluster::Live => user_get_live_bouquet(config, username, &target).await,
XtreamCluster::Video => user_get_vod_bouquet(config, username, &target).await,
XtreamCluster::Series => user_get_series_bouquet(config, username, &target).await,
};
if bouquet.is_none() && category_id.trim().is_empty() {
return None
match bouquet {
None => None,
Some(bouquet_categories) => {
let mut filter = HashSet::new();
let entries: Option<Vec<String>> = if target == TargetType::Xtream {
// xtream filter has PlaylistXtreamCategory
serde_json::from_str::<Vec<PlaylistXtreamCategory>>(&bouquet_categories)
.ok()
.map(|v| v.into_iter().map(|c| c.id).collect())
} else {
// m3u filter has only group names
serde_json::from_str::<Vec<String>>(&bouquet_categories).ok()
};
if let Some(entries) = entries {
filter.extend(entries);
}
Some(filter)
}
}
let mut filter = HashSet::new();
let category_id = category_id.trim();
if !category_id.is_empty() {
filter.insert(category_id.to_string());
}
let bouquet_categories = bouquet.unwrap_or_default();
if bouquet_categories.is_empty() || bouquet_categories == "null" {
return Some(filter).filter(|f| !f.is_empty());
}
let entries: Option<Vec<String>> = if target == TargetType::Xtream {
serde_json::from_str::<Vec<PlaylistXtreamCategory>>(&bouquet_categories)
.ok()
.map(|v| v.into_iter().map(|c| c.id).collect())
} else {
serde_json::from_str::<Vec<String>>(&bouquet_categories).ok()
};
if let Some(entries) = entries {
filter.extend(entries);
}
Some(filter)
}
+3 -4
View File
@@ -25,7 +25,7 @@ impl XtreamPlaylistIterator {
cluster: XtreamCluster,
config: &Config,
target: &ConfigTarget,
category_id: u32,
category_id: Option<u32>,
user: &ProxyUserCredentials,
) -> Result<Self, M3uFilterError> {
if let Some(storage_path) = xtream_get_storage_path(config, target.name.as_str()) {
@@ -41,8 +41,7 @@ impl XtreamPlaylistIterator {
let options = XtreamMappingOptions::from_target_options(target.options.as_ref(), config);
let server_info = config.get_user_server_info(user);
let category_id_filter = if category_id == 0 { String::new() } else { category_id.to_string() };
let filter = user_get_bouquet_filter(config, &user.username, &category_id_filter, TargetType::Xtream, cluster).await;
let filter = user_get_bouquet_filter(config, &user.username, category_id, TargetType::Xtream, cluster).await;
Ok(Self {
reader,
@@ -88,7 +87,7 @@ pub async fn new(
cluster: XtreamCluster,
config: &Config,
target: &ConfigTarget,
category_id: u32,
category_id: Option<u32>,
user: &ProxyUserCredentials,
) -> Result<Self, M3uFilterError> {
Ok(Self {
+1 -1
View File
@@ -430,7 +430,7 @@ pub async fn xtream_load_rewrite_playlist(
cluster: XtreamCluster,
config: &Config,
target: &ConfigTarget,
category_id: u32,
category_id: Option<u32>,
user: &ProxyUserCredentials,
) -> Result<XtreamPlaylistJsonIterator, M3uFilterError> {
XtreamPlaylistJsonIterator::new(cluster, config, target, category_id, user).await