2024-12-10 23:16:04 +01:00
use crate ::m3u_filter_error ::{ M3uFilterError , M3uFilterErrorKind };
use crate ::model ::config ::{ Config , ConfigInput , ConfigTarget , InputType , TargetType };
use crate ::model ::playlist ::{ FetchedPlaylist , PlaylistEntry , PlaylistItem , PlaylistItemType , UUIDType , XtreamCluster };
2024-05-04 20:30:20 +02:00
use crate ::processing ::playlist_processor ::ProcessingPipe ;
2024-12-10 23:16:04 +01:00
use crate ::repository ::storage ::get_input_storage_path ;
2024-12-13 14:28:29 +01:00
use crate ::repository ::xtream_repository ::{ xtream_get_info_file_paths , xtream_update_input_vod_info_file , xtream_update_input_vod_tmdb_file };
2024-12-13 17:51:47 +01:00
use crate ::repository ::IndexedDocumentIndex ;
2024-05-04 20:30:20 +02:00
use crate ::utils ::download ;
2024-12-13 14:28:29 +01:00
use serde_json ::{ Map , Value };
2024-12-10 23:16:04 +01:00
use std ::collections ::HashSet ;
2024-12-12 18:10:47 +01:00
use std ::fs ::File ;
use std ::io ::{ BufWriter , Error , ErrorKind , Write };
2024-12-10 23:16:04 +01:00
use std ::rc ::Rc ;
2024-05-04 20:30:20 +02:00
2024-12-13 14:28:29 +01:00
const TAG_VOD_INFO_INFO : & str = "info" ;
const TAG_VOD_INFO_MOVIE_DATA : & str = "movie_data" ;
const TAG_VOD_INFO_TMDB_ID : & str = "tmdb_id" ;
const TAG_VOD_INFO_STREAM_ID : & str = "stream_id" ;
2024-12-02 17:24:06 +01:00
pub async fn playlist_resolve_series ( target : & ConfigTarget , errors : & mut Vec < M3uFilterError > ,
2024-05-04 20:30:20 +02:00
pipe : & ProcessingPipe ,
2024-12-10 23:16:04 +01:00
provider_fpl : & mut FetchedPlaylist < '_ > ,
processed_fpl : & mut FetchedPlaylist < '_ > ) {
2024-05-04 20:30:20 +02:00
let ( resolve_series , resolve_series_delay ) =
if let Some ( options ) = & target . options {
2024-12-10 23:16:04 +01:00
( options . xtream_resolve_series && provider_fpl . input . input_type == InputType ::Xtream && target . has_output ( & TargetType ::M3u ),
2024-05-04 20:30:20 +02:00
options . xtream_resolve_series_delay )
} else {
( false , 0 )
};
if resolve_series {
2024-12-10 23:16:04 +01:00
// collect all series in the processed lists
let to_process_uuids : HashSet < Rc < UUIDType >> = processed_fpl . playlistgroups . iter ()
. filter ( | plg | plg . xtream_cluster == XtreamCluster ::Series )
. flat_map ( | plg | & plg . channels )
. filter ( | pli | pli . header . borrow (). item_type == PlaylistItemType ::SeriesInfo )
. map ( | pli | Rc ::clone ( & pli . header . borrow (). uuid )). collect ();
let mut series_playlist = download ::get_xtream_playlist_series ( provider_fpl , to_process_uuids , errors , resolve_series_delay ). await ;
2024-05-04 20:30:20 +02:00
// original content saved into original list
for plg in & series_playlist {
2024-12-10 23:16:04 +01:00
provider_fpl . update_playlist ( plg );
2024-05-04 20:30:20 +02:00
}
// run processing pipe over new items
for f in pipe {
let r = f ( & mut series_playlist , target );
if let Some ( v ) = r {
series_playlist = v ;
}
}
// assign new items to the new playlist
for plg in & series_playlist {
2024-12-10 23:16:04 +01:00
processed_fpl . update_playlist ( plg );
2024-05-04 20:30:20 +02:00
}
}
2024-12-10 23:16:04 +01:00
}
2024-12-13 14:28:29 +01:00
async fn playlist_resolve_vod_process_playlist_item ( pli : & PlaylistItem , input : & ConfigInput , errors : & mut Vec < M3uFilterError > , resolve_delay : u16 ) -> Option < String > {
2024-12-10 23:16:04 +01:00
let mut result = None ;
let provider_id = pli . get_provider_id (). unwrap_or ( 0 );
if let Some ( info_url ) = download ::get_xtream_player_api_info_url ( input , XtreamCluster ::Video , provider_id ) {
2024-12-13 00:13:14 +01:00
result = match download ::get_xtream_stream_info_content ( & info_url , input ). await {
2024-12-10 23:16:04 +01:00
Ok ( content ) => Some ( content ),
Err ( err ) => {
errors . push ( M3uFilterError ::new ( M3uFilterErrorKind ::Info , format! ( " {err} " )));
None
}
};
}
if resolve_delay > 0 {
actix_web ::rt ::time ::sleep ( std ::time ::Duration ::new ( u64 ::from ( resolve_delay ), 0 )). await ;
}
result
}
2024-12-13 14:28:29 +01:00
fn write_vod_info_content_to_temp_file ( writer : & mut BufWriter <& File > , provider_id : u32 , content : & str ) -> std ::io ::Result < () > {
2024-12-12 18:10:47 +01:00
let length = u32 ::try_from ( content . len ()). map_err ( | err | Error ::new ( ErrorKind ::Other , err )) ? ;
if length > 0 {
2024-12-13 14:28:29 +01:00
writer . write_all ( & provider_id . to_le_bytes ()) ? ;
2024-12-12 18:10:47 +01:00
writer . write_all ( & length . to_le_bytes ()) ? ;
writer . write_all ( content . as_bytes ()) ? ;
}
Ok (())
}
2024-12-13 14:28:29 +01:00
fn write_vod_info_tmdb_to_temp_file ( writer : & mut BufWriter <& File > , provider_id : u32 , tmdb_id : u32 ) -> std ::io ::Result < () > {
writer . write_all ( & provider_id . to_le_bytes ()) ? ;
writer . write_all ( & tmdb_id . to_le_bytes ()) ? ;
Ok (())
}
fn get_resolve_video_options ( target : & ConfigTarget , fpl : & FetchedPlaylist ) -> ( bool , u16 ) {
2024-12-10 23:16:04 +01:00
let ( resolve_movies , resolve_delay ) =
2024-12-13 14:28:29 +01:00
target . options . as_ref (). map_or (( false , 0 ), | opt | ( opt . xtream_resolve_video && fpl . input . input_type == InputType ::Xtream , opt . xtream_resolve_video_delay ));
2024-12-12 18:10:47 +01:00
( resolve_movies , resolve_delay )
}
2024-12-13 14:28:29 +01:00
fn get_u32_from_serde_value ( value : & Value ) -> Option < u32 > {
match value {
Value ::Number ( num_val ) => num_val . as_u64 (). and_then ( | val | u32 ::try_from ( val ). ok ()),
Value ::String ( str_val ) => {
match str_val . parse ::< u32 > () {
Ok ( sid ) => Some ( sid ),
Err ( _ ) => None
}
}
_ => None ,
2024-12-12 18:10:47 +01:00
}
}
2024-12-13 14:28:29 +01:00
fn extract_provider_id_and_tmdb_id_from_vod_info ( content : & str ) -> Option < ( u32 , u32 ) > {
if let Ok ( mut doc ) = serde_json ::from_str ::< Map < String , Value >> ( content ) {
if let Some ( Value ::Object ( movie_data )) = doc . get_mut ( TAG_VOD_INFO_MOVIE_DATA ) {
if let Some ( stream_id_value ) = movie_data . get ( TAG_VOD_INFO_STREAM_ID ) {
if let Some ( stream_id ) = get_u32_from_serde_value ( stream_id_value ) {
if let Some ( Value ::Object ( info )) = doc . get_mut ( TAG_VOD_INFO_INFO ) {
if let Some ( tmdb_id_value ) = info . get ( TAG_VOD_INFO_TMDB_ID ) {
if let Some ( tmdb_id ) = get_u32_from_serde_value ( tmdb_id_value ) {
return Some (( stream_id , tmdb_id ));
}
}
}
return Some (( stream_id , 0 ));
}
}
}
2024-12-10 23:16:04 +01:00
}
2024-12-13 14:28:29 +01:00
None
}
2024-12-12 18:10:47 +01:00
2024-12-13 14:28:29 +01:00
fn create_resolve_vod_info_temp_files ( errors : & mut Vec < M3uFilterError > ) -> Option < ( File , File ) > {
let temp_file_info = match tempfile ::tempfile () {
2024-12-12 18:10:47 +01:00
Ok ( value ) => value ,
Err ( err ) => {
2024-12-13 14:28:29 +01:00
errors . push ( M3uFilterError ::new ( M3uFilterErrorKind ::Info , format! ( "Cant resolve vod, could not create temporary file {err} " )));
return None ;
2024-12-12 18:10:47 +01:00
}
};
2024-12-13 14:28:29 +01:00
let temp_file_tmdb = match tempfile ::tempfile () {
Ok ( value ) => value ,
Err ( err ) => {
errors . push ( M3uFilterError ::new ( M3uFilterErrorKind ::Info , format! ( "Cant resolve vod tmdb, could not create temporary file {err} " )));
return None ;
}
};
Some (( temp_file_info , temp_file_tmdb ))
}
pub async fn playlist_resolve_vod ( cfg : & Config , target : & ConfigTarget , errors : & mut Vec < M3uFilterError > , fpl : & FetchedPlaylist < '_ > ) {
let ( resolve_movies , resolve_delay ) = get_resolve_video_options ( target , fpl );
if ! resolve_movies { return ; }
// we cant write to the indexed-document directly because of the write lock and time-consuming operation.
// All readers would be waiting for the lock and the app would be unresponsive.
// We collect the content into a temp file and write it once we collected everything.
2024-12-13 17:51:47 +01:00
let Some (( mut temp_file_info , mut temp_file_tmdb )) = create_resolve_vod_info_temp_files ( errors ) else { return };
2024-12-12 18:10:47 +01:00
let mut processed_vod_ids = read_processed_vod_info_ids ( cfg , errors , fpl ). await ;
2024-12-13 14:28:29 +01:00
let mut info_writer = BufWriter ::new ( & temp_file_info );
let mut tmdb_writer = BufWriter ::new ( & temp_file_tmdb );
2024-12-13 17:51:47 +01:00
let mut info_updated = false ;
let mut tmdb_updated = false ;
2024-12-12 18:10:47 +01:00
for pli in fpl . playlistgroups . iter (). flat_map ( | plg | & plg . channels ) {
2024-12-13 14:28:29 +01:00
let a = pli . header . borrow_mut (). get_provider_id (). as_ref (). map_or ( false , | pid | processed_vod_ids . contains ( pid ));
if ! a {
if let Some ( content ) = playlist_resolve_vod_process_playlist_item ( pli , fpl . input , errors , resolve_delay ). await {
if let Some (( provider_id , tmdb_id )) = extract_provider_id_and_tmdb_id_from_vod_info ( & content ) {
if let Err ( err ) = write_vod_info_content_to_temp_file ( & mut info_writer , provider_id , & content ) {
errors . push ( M3uFilterError ::new ( M3uFilterErrorKind ::Notify , format! ( "Failed to resolve vod, could not write to temporary file {err} " )));
return ;
}
2024-12-13 17:51:47 +01:00
info_updated = true ;
2024-12-13 14:28:29 +01:00
processed_vod_ids . insert ( provider_id );
if tmdb_id > 0 {
if let Err ( err ) = write_vod_info_tmdb_to_temp_file ( & mut tmdb_writer , provider_id , tmdb_id ) {
errors . push ( M3uFilterError ::new ( M3uFilterErrorKind ::Notify , format! ( "Failed to resolve vod tmdb, could not write to temporary file {err} " )));
return ;
}
2024-12-13 17:51:47 +01:00
tmdb_updated = true ;
2024-12-13 14:28:29 +01:00
}
2024-12-12 18:10:47 +01:00
}
}
}
}
2024-12-13 17:51:47 +01:00
if info_updated {
if let Err ( err ) = info_writer . flush () {
errors . push ( M3uFilterError ::new ( M3uFilterErrorKind ::Notify , format! ( "Failed to resolve vod, could not write to temporary file {err} " )));
}
drop ( info_writer );
if let Err ( err ) = xtream_update_input_vod_info_file ( cfg , fpl . input , & mut temp_file_info ). await {
errors . push ( err );
}
2024-12-12 18:10:47 +01:00
}
2024-12-13 17:51:47 +01:00
if tmdb_updated {
if let Err ( err ) = tmdb_writer . flush () {
errors . push ( M3uFilterError ::new ( M3uFilterErrorKind ::Notify , format! ( "Failed to resolve vod tmdb, could not write to temporary file {err} " )));
}
drop ( tmdb_writer );
if let Err ( err ) = xtream_update_input_vod_tmdb_file ( cfg , fpl . input , & mut temp_file_tmdb ). await {
errors . push ( err );
}
2024-12-13 14:28:29 +01:00
}
2024-12-12 18:10:47 +01:00
}
2024-12-13 14:28:29 +01:00
async fn read_processed_vod_info_ids ( cfg : & Config , errors : & mut Vec < M3uFilterError > , fpl : & FetchedPlaylist < '_ > ) -> HashSet < u32 > {
2024-12-12 18:10:47 +01:00
let mut processed_vod_ids = HashSet ::new ();
{
match get_input_storage_path ( fpl . input , & cfg . working_dir ). map ( | storage_path | xtream_get_info_file_paths ( & storage_path , XtreamCluster ::Video )) {
Ok ( Some (( file_path , idx_path ))) => {
match cfg . file_locks . read_lock ( & file_path ). await {
Ok ( file_lock ) => {
2024-12-13 17:51:47 +01:00
if let Ok ( info_id_mapping ) = IndexedDocumentIndex ::< u32 > ::load ( & idx_path ) {
2024-12-12 18:10:47 +01:00
info_id_mapping . traverse ( | keys , _ | {
2024-12-13 14:28:29 +01:00
for doc_id in keys { processed_vod_ids . insert ( * doc_id ); }
2024-12-12 18:10:47 +01:00
});
2024-12-10 23:16:04 +01:00
}
2024-12-12 18:10:47 +01:00
drop ( file_lock );
2024-12-10 23:16:04 +01:00
}
2024-12-12 18:10:47 +01:00
Err ( err ) => errors . push ( M3uFilterError ::new ( M3uFilterErrorKind ::Info , format! ( " {err} " ))),
2024-12-10 23:16:04 +01:00
}
}
2024-12-12 18:10:47 +01:00
Ok ( None ) => errors . push ( M3uFilterError ::new ( M3uFilterErrorKind ::Notify , format! ( "Could not create storage path for input {} " , & fpl . input . name . as_ref (). map_or ( "?" , | v | v )))),
Err ( err ) => errors . push ( M3uFilterError ::new ( M3uFilterErrorKind ::Notify , format! ( "Could not create storage path for input {err} " ))),
2024-12-10 23:16:04 +01:00
}
}
2024-12-12 18:10:47 +01:00
processed_vod_ids
2024-12-13 14:28:29 +01:00
}