mirror of
https://github.com/euzu/tuliprox.git
synced 2026-10-08 00:42:32 +02:00
Feature/bouquet editor (#854)
New Features
Added target-specific bouquet filtering with whitelist and blacklist modes for Live, VOD, and Series groups.
Added an in-context bouquet editor with search, selection controls, status indicators, reset, and streaming previews.
Targets are managed by name, with saved selections available through the Source Editor.
Improved stream alias resolution and provider failover behavior.
Bug Fixes
Existing playlists are preserved when refreshes produce no usable items.
Empty playlists are no longer published.
Configuration changes now recover more safely from persistence failures.
Documentation
Updated feature and configuration documentation and added translations.
This commit is contained in:
@@ -258,6 +258,15 @@ pub(crate) async fn prepare_playlist_for_target<E: EventSink + Clone + 'static,
|
||||
debug!("Executing processing pipes");
|
||||
let broadcast_step = create_broadcast_callback(&ctx.events);
|
||||
|
||||
let bouquet_file =
|
||||
tuliprox_repository::load_target_bouquet(&ctx.config, &target.name).await.map_err(|err| vec![err])?;
|
||||
let bouquet_filter =
|
||||
bouquet_file.and_then(|file| tuliprox_core::model::TargetBouquetFilter::from_dto(file.bouquet));
|
||||
if let Some(ref filter) = bouquet_filter {
|
||||
let (live, vod, series) = filter.cluster_counts();
|
||||
debug!("Loaded target bouquet for '{}': live={:?}, vod={:?}, series={:?}", target.name, live, vod, series);
|
||||
}
|
||||
|
||||
let pipe = get_processing_pipe(target);
|
||||
let mut step = StepMeasure::new(&target.name, broadcast_step);
|
||||
for provider_fpl in playlists.iter_mut() {
|
||||
@@ -266,7 +275,7 @@ pub(crate) async fn prepare_playlist_for_target<E: EventSink + Clone + 'static,
|
||||
);
|
||||
step.broadcast("Executing transformations on '{}' playlist", &target.name);
|
||||
let (mut processed_fpl, input_outcome) =
|
||||
execute_pipe(target, &pipe, provider_fpl, &mut duplicates, consume_input_source)
|
||||
execute_pipe(target, &pipe, provider_fpl, &mut duplicates, consume_input_source, bouquet_filter.as_ref())
|
||||
.map_err(|err| vec![err])?;
|
||||
debug!("Target '{}' input '{}' pipeline outcome: {input_outcome:?}", target.name, provider_fpl.input.name);
|
||||
aggregate_outcome.merge(input_outcome);
|
||||
@@ -423,10 +432,6 @@ pub(crate) async fn finalize_prepared_target<E: EventSink + Clone + 'static, M:
|
||||
apply_persist_filter(target, &mut flat_new_playlist);
|
||||
retain_epg_referenced_by_groups(&flat_new_playlist, &mut new_epg);
|
||||
|
||||
if process_watch(&ctx.config, &ctx.events, target, &flat_new_playlist).await {
|
||||
step.tick("group watches");
|
||||
log_memory_snapshot(format!("target '{}' after_group_watches", target.name).as_str());
|
||||
}
|
||||
let merged_epg = if ctx.config.config.load().disk_based_processing {
|
||||
// Per-source drain to disk, then multi-way merge. Errors are pushed
|
||||
// to `errors` rather than `?` because the function returns
|
||||
@@ -458,6 +463,10 @@ pub(crate) async fn finalize_prepared_target<E: EventSink + Clone + 'static, M:
|
||||
ctx.playlist_state.as_ref(),
|
||||
)
|
||||
.await;
|
||||
if result.is_ok() && process_watch(&ctx.config, &ctx.events, target, &flat_new_playlist).await {
|
||||
step.tick("group watches");
|
||||
log_memory_snapshot(format!("target '{}' after_group_watches", target.name).as_str());
|
||||
}
|
||||
step.stop("Persisting playlists");
|
||||
log_memory_snapshot(format!("target '{}' after_persist", target.name).as_str());
|
||||
(result, errors)
|
||||
|
||||
@@ -160,8 +160,8 @@ fn execute_pipe_freezes_input_stream_id_without_rename_or_mapper() {
|
||||
let mut duplicates = HashSet::new();
|
||||
let target = ConfigTarget::from(&ConfigTargetDto::default());
|
||||
|
||||
let (mut processed, _outcome) =
|
||||
execute_pipe(&target, &vec![], &mut fetched, &mut duplicates, false).expect("target processing should succeed");
|
||||
let (mut processed, _outcome) = execute_pipe(&target, &vec![], &mut fetched, &mut duplicates, false, None)
|
||||
.expect("target processing should succeed");
|
||||
let mut groups = processed.source.take_groups();
|
||||
|
||||
assert_eq!(groups[0].channels[0].header.input_stream_id.as_ref(), "origin-alpha");
|
||||
@@ -169,6 +169,111 @@ fn execute_pipe_freezes_input_stream_id_without_rename_or_mapper() {
|
||||
assert_eq!(groups[0].channels[0].header.input_stream_id.as_ref(), "origin-alpha");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn execute_pipe_applies_target_bouquet_prefilter() {
|
||||
let input = ConfigInput::default();
|
||||
let item1 = PlaylistItem {
|
||||
header: PlaylistItemHeader {
|
||||
id: "ch-1".intern(),
|
||||
group: "Kids".intern(),
|
||||
xtream_cluster: XtreamCluster::Live,
|
||||
..Default::default()
|
||||
},
|
||||
};
|
||||
let item2 = PlaylistItem {
|
||||
header: PlaylistItemHeader {
|
||||
id: "ch-2".intern(),
|
||||
group: "Adults".intern(),
|
||||
xtream_cluster: XtreamCluster::Live,
|
||||
..Default::default()
|
||||
},
|
||||
};
|
||||
let source = MemoryPlaylistSource::new(vec![
|
||||
PlaylistGroup { id: 1, title: "Kids".intern(), channels: vec![item1], xtream_cluster: XtreamCluster::Live },
|
||||
PlaylistGroup { id: 2, title: "Adults".intern(), channels: vec![item2], xtream_cluster: XtreamCluster::Live },
|
||||
])
|
||||
.into_source();
|
||||
let mut fetched = FetchedPlaylist { input: &input, source, epg: None };
|
||||
let mut duplicates = HashSet::new();
|
||||
let target = ConfigTarget::from(&ConfigTargetDto::default());
|
||||
|
||||
let bouquet_dto =
|
||||
shared::model::PlaylistClusterBouquetDto { live: Some(vec!["Kids".to_string()]), vod: None, series: None };
|
||||
let filter =
|
||||
tuliprox_core::model::TargetBouquetFilter::from_dto(shared::model::TargetBouquetDto::whitelist(bouquet_dto))
|
||||
.unwrap();
|
||||
|
||||
let (mut processed, _outcome) = execute_pipe(&target, &vec![], &mut fetched, &mut duplicates, false, Some(&filter))
|
||||
.expect("target processing should succeed");
|
||||
let groups = processed.source.take_groups();
|
||||
|
||||
assert_eq!(groups.len(), 1);
|
||||
assert_eq!(groups[0].title.as_ref(), "Kids");
|
||||
assert_eq!(groups[0].channels.len(), 1);
|
||||
assert_eq!(groups[0].channels[0].header.id.as_ref(), "ch-1");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn execute_pipe_prefilter_does_not_suppress_allowed_duplicate() {
|
||||
let input = ConfigInput::default();
|
||||
// Two items with the same URL (same UUID): one in disallowed group, one in allowed group.
|
||||
let item_disallowed = PlaylistItem {
|
||||
header: PlaylistItemHeader {
|
||||
id: "ch-1".intern(),
|
||||
url: "http://provider.example/stream.ts".intern(),
|
||||
group: "Adults".intern(),
|
||||
xtream_cluster: XtreamCluster::Live,
|
||||
..Default::default()
|
||||
},
|
||||
};
|
||||
let item_allowed = PlaylistItem {
|
||||
header: PlaylistItemHeader {
|
||||
id: "ch-2".intern(),
|
||||
url: "http://provider.example/stream.ts".intern(),
|
||||
group: "Kids".intern(),
|
||||
xtream_cluster: XtreamCluster::Live,
|
||||
..Default::default()
|
||||
},
|
||||
};
|
||||
let source = MemoryPlaylistSource::new(vec![
|
||||
PlaylistGroup {
|
||||
id: 1,
|
||||
title: "Adults".intern(),
|
||||
channels: vec![item_disallowed],
|
||||
xtream_cluster: XtreamCluster::Live,
|
||||
},
|
||||
PlaylistGroup {
|
||||
id: 2,
|
||||
title: "Kids".intern(),
|
||||
channels: vec![item_allowed],
|
||||
xtream_cluster: XtreamCluster::Live,
|
||||
},
|
||||
])
|
||||
.into_source();
|
||||
let mut fetched = FetchedPlaylist { input: &input, source, epg: None };
|
||||
let mut duplicates = HashSet::new();
|
||||
let target = ConfigTarget::from(&ConfigTargetDto {
|
||||
options: Some(ConfigTargetOptions { remove_duplicates: true, ..Default::default() }),
|
||||
..Default::default()
|
||||
});
|
||||
|
||||
let bouquet_dto =
|
||||
shared::model::PlaylistClusterBouquetDto { live: Some(vec!["Kids".to_string()]), vod: None, series: None };
|
||||
let filter =
|
||||
tuliprox_core::model::TargetBouquetFilter::from_dto(shared::model::TargetBouquetDto::whitelist(bouquet_dto))
|
||||
.unwrap();
|
||||
|
||||
let (mut processed, _outcome) = execute_pipe(&target, &vec![], &mut fetched, &mut duplicates, false, Some(&filter))
|
||||
.expect("target processing should succeed");
|
||||
let groups = processed.source.take_groups();
|
||||
|
||||
// The allowed item must be retained because the rejected item did not consume the duplicate UUID slot.
|
||||
assert_eq!(groups.len(), 1);
|
||||
assert_eq!(groups[0].title.as_ref(), "Kids");
|
||||
assert_eq!(groups[0].channels.len(), 1);
|
||||
assert_eq!(groups[0].channels[0].header.id.as_ref(), "ch-2");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn legacy_messagepack_playlist_items_default_missing_input_stream_id() {
|
||||
let mut source = PlaylistItem {
|
||||
@@ -1330,7 +1435,7 @@ match {
|
||||
FetchedPlaylist { input: &input, source: memory_source(vec![channel.clone(), channel]), epg: None };
|
||||
let mut duplicates = HashSet::new();
|
||||
let (mut processed, _outcome) =
|
||||
execute_pipe(&target, &get_processing_pipe(&target), &mut fetched, &mut duplicates, false)
|
||||
execute_pipe(&target, &get_processing_pipe(&target), &mut fetched, &mut duplicates, false, None)
|
||||
.expect("processing pipe must run");
|
||||
assert_eq!(processed.get_channel_count(), 1, "processing pipe must remove the duplicate");
|
||||
|
||||
|
||||
@@ -1,5 +1,6 @@
|
||||
#![allow(clippy::wildcard_imports)]
|
||||
use super::*;
|
||||
use shared::model::PlaylistEntry;
|
||||
|
||||
pub(crate) fn retain_playlist_items(
|
||||
source: &mut PlaylistSource,
|
||||
@@ -478,8 +479,9 @@ pub(crate) fn execute_pipe<'a>(
|
||||
fpl: &mut FetchedPlaylist<'a>,
|
||||
duplicates: &mut HashSet<UUIDType>,
|
||||
consume_source: bool,
|
||||
bouquet_filter: Option<&tuliprox_core::model::TargetBouquetFilter>,
|
||||
) -> Result<(FetchedPlaylist<'a>, PipelineOutcome), TuliproxError> {
|
||||
let source = if consume_source {
|
||||
let mut source = if consume_source {
|
||||
if fpl.is_memory() {
|
||||
MemoryPlaylistSource::new(fpl.source.take_groups()).into_source()
|
||||
} else {
|
||||
@@ -489,21 +491,58 @@ pub(crate) fn execute_pipe<'a>(
|
||||
fpl.clone_source()?
|
||||
};
|
||||
|
||||
let mut new_fpl = FetchedPlaylist { input: fpl.input, source, epg: fpl.epg.clone() };
|
||||
// In-memory items are frozen here at the target-processing boundary. Read-only disk sources
|
||||
// capture the same identity when their persisted M3U/Xtream items are converted to PlaylistItem.
|
||||
if new_fpl.is_memory() {
|
||||
for item in new_fpl.items_mut() {
|
||||
let is_memory = source.is_memory();
|
||||
let pre_transform_dedup_enabled = target.execution_plan.pre_transform_identity_dedup;
|
||||
if !is_memory && pre_transform_dedup_enabled {
|
||||
warn!(
|
||||
"Target '{}' input '{}': pre_transform_identity_dedup has no effect for disk-based playlist sources",
|
||||
target.name, fpl.input.name
|
||||
);
|
||||
}
|
||||
// Memory sources dedup here; disk sources cannot, because persisted B+Tree
|
||||
// leaves are read-only at the pipeline boundary. Cross-input dedup for disk
|
||||
// sources runs later in `deduplicate_playlist` during target finalization.
|
||||
let deduplicate_memory = is_memory && pre_transform_dedup_enabled;
|
||||
let mut items = Vec::new();
|
||||
let mut rejected = 0usize;
|
||||
let mut duplicates_removed = 0usize;
|
||||
let mut total_items = 0usize;
|
||||
|
||||
for mut item in source.into_items() {
|
||||
total_items += 1;
|
||||
if is_memory {
|
||||
item.header.freeze_input_stream_id();
|
||||
}
|
||||
}
|
||||
if target.execution_plan.pre_transform_identity_dedup {
|
||||
new_fpl.deduplicate(duplicates);
|
||||
if bouquet_filter.is_some_and(|filter| !filter.allows(item.header.xtream_cluster, &item.header.group)) {
|
||||
rejected += 1;
|
||||
continue;
|
||||
}
|
||||
if deduplicate_memory && !duplicates.insert(item.get_uuid()) {
|
||||
duplicates_removed += 1;
|
||||
continue;
|
||||
}
|
||||
items.push(item);
|
||||
}
|
||||
|
||||
if rejected > 0 {
|
||||
debug!(
|
||||
"Target '{}' bouquet prefilter for input '{}': rejected {}/{} items",
|
||||
target.name, fpl.input.name, rejected, total_items
|
||||
);
|
||||
}
|
||||
if duplicates_removed > 0 {
|
||||
debug!(
|
||||
"Target '{}' input '{}': pre-transform dedup removed {} duplicate items",
|
||||
target.name, fpl.input.name, duplicates_removed
|
||||
);
|
||||
}
|
||||
|
||||
let items = new_fpl.source.into_items().collect();
|
||||
let (groups, outcome) = execute_pipeline_on_items(items, target, pipe);
|
||||
new_fpl.source = MemoryPlaylistSource::new(groups).into_source();
|
||||
let new_fpl = FetchedPlaylist {
|
||||
input: fpl.input,
|
||||
source: MemoryPlaylistSource::new(groups).into_source(),
|
||||
epg: fpl.epg.clone(),
|
||||
};
|
||||
Ok((new_fpl, outcome))
|
||||
}
|
||||
|
||||
|
||||
@@ -44,6 +44,10 @@ pub enum StalkerCluster {
|
||||
Series,
|
||||
}
|
||||
|
||||
fn raw_group_catalog_storage_path(stalker_storage_path: &std::path::Path) -> &std::path::Path {
|
||||
stalker_storage_path.parent().unwrap_or(stalker_storage_path)
|
||||
}
|
||||
|
||||
const DEFAULT_STALKER_CLUSTERS: [StalkerCluster; 3] =
|
||||
[StalkerCluster::Live, StalkerCluster::Vod, StalkerCluster::Series];
|
||||
|
||||
@@ -168,6 +172,7 @@ pub async fn download_stalker_playlist(
|
||||
);
|
||||
}
|
||||
};
|
||||
let catalog_storage_path = raw_group_catalog_storage_path(&storage_path);
|
||||
|
||||
let outcome = if let Some(_refresh_permit) = try_acquire_stalker_refresh(&input.name).await {
|
||||
let handshake = match api_client.handshake().await {
|
||||
@@ -297,6 +302,22 @@ pub async fn download_stalker_playlist(
|
||||
Ok(items) => {
|
||||
counts[cluster as usize] = items.len();
|
||||
let cluster_groups = groups_for_cluster(items, cluster, &input.name);
|
||||
let group_titles = cluster_groups.iter().map(|g| g.title.to_string()).collect::<Vec<String>>();
|
||||
let xc = xtream_cluster(cluster);
|
||||
if let Err(publish_err) = tuliprox_repository::publish_raw_group_catalog(
|
||||
catalog_storage_path,
|
||||
&input.name,
|
||||
xc,
|
||||
group_titles,
|
||||
&app_config.file_locks,
|
||||
)
|
||||
.await
|
||||
{
|
||||
warn!(
|
||||
"Failed to publish raw group catalog for stalker input '{}' cluster {xc:?}: {publish_err}",
|
||||
input.name
|
||||
);
|
||||
}
|
||||
groups.extend(cluster_groups);
|
||||
}
|
||||
Err(err) => errors.push(err),
|
||||
@@ -339,12 +360,7 @@ fn groups_for_cluster(
|
||||
cluster: StalkerCluster,
|
||||
input_name: &str,
|
||||
) -> Vec<PlaylistGroup> {
|
||||
use shared::model::XtreamCluster;
|
||||
let xtream_cluster = match cluster {
|
||||
StalkerCluster::Live => XtreamCluster::Live,
|
||||
StalkerCluster::Vod => XtreamCluster::Video,
|
||||
StalkerCluster::Series => XtreamCluster::Series,
|
||||
};
|
||||
let xtream_cluster = xtream_cluster(cluster);
|
||||
let mut groups_map: indexmap::IndexMap<u32, PlaylistGroup> = indexmap::IndexMap::new();
|
||||
for item in items {
|
||||
let category_id = item.category_id;
|
||||
@@ -360,6 +376,14 @@ fn groups_for_cluster(
|
||||
groups_map.into_values().collect()
|
||||
}
|
||||
|
||||
const fn xtream_cluster(cluster: StalkerCluster) -> shared::model::XtreamCluster {
|
||||
match cluster {
|
||||
StalkerCluster::Live => shared::model::XtreamCluster::Live,
|
||||
StalkerCluster::Vod => shared::model::XtreamCluster::Video,
|
||||
StalkerCluster::Series => shared::model::XtreamCluster::Series,
|
||||
}
|
||||
}
|
||||
|
||||
/// Map a Stalker failure onto the workspace error type, keeping its classification.
|
||||
///
|
||||
/// Flattening everything to `ProviderConnection` read as "the network had a bad moment",
|
||||
@@ -478,6 +502,12 @@ mod tests {
|
||||
|
||||
fn runtime_cfg() -> StalkerInputConfig { StalkerInputConfig::default() }
|
||||
|
||||
#[test]
|
||||
fn raw_group_catalog_is_published_at_the_input_storage_root() {
|
||||
let stalker_path = std::path::Path::new("/storage/input_portal/stalker");
|
||||
assert_eq!(raw_group_catalog_storage_path(stalker_path), std::path::Path::new("/storage/input_portal"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn runtime_client_cache_key_changes_with_endpoint_preference() {
|
||||
let mut cfg = runtime_cfg();
|
||||
|
||||
@@ -19,8 +19,8 @@ use tuliprox_iptv::stalker::{
|
||||
};
|
||||
use tuliprox_repository::{
|
||||
stalker_generation_repository::{
|
||||
cleanup_obsolete_generations, clear_checkpoint, generation_data_path, load_checkpoint, publish_selection,
|
||||
save_checkpoint, StalkerCheckpoint, StalkerGenerationData, StalkerRefreshPhase,
|
||||
cleanup_obsolete_generations, clear_checkpoint, generation_data_path, load_active_manifest, load_checkpoint,
|
||||
publish_selection, save_checkpoint, StalkerCheckpoint, StalkerGenerationData, StalkerRefreshPhase,
|
||||
},
|
||||
stalker_repository::{
|
||||
load_stalker_items_after, prepare_stalker_episode_series_at, promote_stalker_file, remove_stalker_file,
|
||||
@@ -183,6 +183,42 @@ async fn load_or_start_checkpoint(
|
||||
Ok(state)
|
||||
}
|
||||
|
||||
async fn finish_completed_refresh_checked(
|
||||
app_config: &Arc<AppConfig>,
|
||||
storage_path: &Path,
|
||||
identity_fingerprint: u64,
|
||||
checkpoint: &StalkerCheckpoint,
|
||||
) -> Result<(), TuliproxError> {
|
||||
if checkpoint.processed == 0 {
|
||||
let active = load_active_manifest(storage_path, identity_fingerprint).await?;
|
||||
let mut retained_has_items = false;
|
||||
if checkpoint.selection_mask & 0b0001 == 0 {
|
||||
if let Some(files) = active.live.as_ref() {
|
||||
retained_has_items = !load_stalker_items_after(app_config, &files.data, None, 1).await?.is_empty();
|
||||
}
|
||||
}
|
||||
if !retained_has_items && checkpoint.selection_mask & 0b0010 == 0 {
|
||||
if let Some(files) = active.vod.as_ref() {
|
||||
retained_has_items = !load_stalker_items_after(app_config, &files.data, None, 1).await?.is_empty();
|
||||
}
|
||||
}
|
||||
if !retained_has_items && checkpoint.selection_mask & 0b0100 == 0 {
|
||||
if let Some(files) = active.series.as_ref() {
|
||||
retained_has_items = !load_stalker_items_after(app_config, &files.roots, None, 1).await?.is_empty()
|
||||
|| !load_stalker_items_after(app_config, &files.episodes, None, 1).await?.is_empty();
|
||||
}
|
||||
}
|
||||
if !retained_has_items {
|
||||
clear_checkpoint(storage_path).await?;
|
||||
cleanup_obsolete_generations(storage_path, &active).await?;
|
||||
return Err(TuliproxError::RepositoryPlaylist(
|
||||
"Refusing to publish an empty Stalker playlist; existing data was retained".to_string(),
|
||||
));
|
||||
}
|
||||
}
|
||||
finish_completed_refresh(storage_path, identity_fingerprint, checkpoint).await
|
||||
}
|
||||
|
||||
async fn finish_completed_refresh(
|
||||
storage_path: &Path,
|
||||
identity_fingerprint: u64,
|
||||
@@ -532,7 +568,7 @@ pub async fn advance_stalker_refresh(
|
||||
checkpoint.retry_count = 0;
|
||||
}
|
||||
StalkerRefreshPhase::Complete => {
|
||||
finish_completed_refresh(storage_path, identity_fingerprint, &checkpoint).await?;
|
||||
finish_completed_refresh_checked(app_config, storage_path, identity_fingerprint, &checkpoint).await?;
|
||||
return Ok(StalkerRefreshOutcome::Complete);
|
||||
}
|
||||
StalkerRefreshPhase::Terminal => {
|
||||
@@ -547,6 +583,42 @@ pub async fn advance_stalker_refresh(
|
||||
#[cfg(test)]
|
||||
mod tests {
|
||||
use super::*;
|
||||
use arc_swap::{ArcSwap, ArcSwapOption};
|
||||
use shared::model::ConfigPaths;
|
||||
use tuliprox_core::{
|
||||
model::{ApiProxyConfig, Config, CustomStreamResponse, HdHomeRunConfig, MediaToolCapabilities, SourcesConfig},
|
||||
utils::FileLockManager,
|
||||
};
|
||||
|
||||
fn test_app_config(storage_dir: &Path) -> Arc<AppConfig> {
|
||||
Arc::new(AppConfig {
|
||||
config: Arc::new(ArcSwap::from_pointee(Config {
|
||||
storage_dir: storage_dir.to_string_lossy().into_owned(),
|
||||
..Config::default()
|
||||
})),
|
||||
sources: Arc::new(ArcSwap::from_pointee(SourcesConfig::default())),
|
||||
hdhomerun: Arc::new(ArcSwapOption::<HdHomeRunConfig>::default()),
|
||||
api_proxy: Arc::new(ArcSwapOption::<ApiProxyConfig>::default()),
|
||||
file_locks: Arc::new(FileLockManager::default()),
|
||||
paths: Arc::new(ArcSwap::from_pointee(ConfigPaths {
|
||||
home_path: String::new(),
|
||||
config_path: String::new(),
|
||||
storage_path: storage_dir.to_string_lossy().into_owned(),
|
||||
config_file_path: String::new(),
|
||||
sources_file_path: String::new(),
|
||||
mapping_file_path: None,
|
||||
mapping_files_used: None,
|
||||
template_file_path: None,
|
||||
template_files_used: None,
|
||||
api_proxy_file_path: String::new(),
|
||||
custom_stream_response_path: None,
|
||||
})),
|
||||
custom_stream_response: Arc::new(ArcSwapOption::<CustomStreamResponse>::default()),
|
||||
access_token_secret: [0; 32],
|
||||
encrypt_secret: [0; 16],
|
||||
media_tools: Arc::new(MediaToolCapabilities::new()),
|
||||
})
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn terminal_checkpoint_is_cleared_after_restart() -> Result<(), Box<dyn std::error::Error>> {
|
||||
@@ -578,6 +650,30 @@ mod tests {
|
||||
Ok(())
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn empty_complete_refresh_keeps_the_active_manifest() -> Result<(), Box<dyn std::error::Error>> {
|
||||
let temp = tempfile::tempdir()?;
|
||||
let app_config = test_app_config(temp.path());
|
||||
let mut active = tuliprox_repository::stalker_generation_repository::StalkerActiveManifest::empty(17);
|
||||
active.live = Some(tuliprox_repository::stalker_generation_repository::ClusterFiles {
|
||||
generation: 11,
|
||||
data: temp.path().join("old-live.db"),
|
||||
});
|
||||
tuliprox_repository::stalker_generation_repository::save_active_manifest(temp.path(), &active).await?;
|
||||
|
||||
let selection = StalkerClusterSelection { live: true, vod: false, series: false, epg: false };
|
||||
let mut checkpoint = StalkerCheckpoint::new(17, 23, selection.mask(), 123);
|
||||
checkpoint.phase = StalkerRefreshPhase::Complete;
|
||||
save_checkpoint(temp.path(), &checkpoint).await?;
|
||||
|
||||
let result = finish_completed_refresh_checked(&app_config, temp.path(), 17, &checkpoint).await;
|
||||
|
||||
assert!(result.is_err());
|
||||
assert_eq!(load_active_manifest(temp.path(), 17).await?, active);
|
||||
assert!(load_checkpoint(temp.path(), 17).await?.is_none());
|
||||
Ok(())
|
||||
}
|
||||
|
||||
#[tokio::test]
|
||||
async fn complete_publication_is_idempotent_after_restart() -> Result<(), Box<dyn std::error::Error>> {
|
||||
let temp = tempfile::tempdir()?;
|
||||
|
||||
Reference in New Issue
Block a user