feat: сортировка заданий по дате и времени, последовательная отправка
- Добавлена сортировка заданий по дате и времени события (самые ранние — первыми) - Добавлена сортировка по варианту (оригинал -> today -> tomorrow) - Изменена отправка на последовательную для гарантированного порядка в очереди - Исправлена нормализация дефисов (2+ дефиса -> 1) - Исправлен SINGLE шаблон (убран слой TEAMS) - Добавлена поддержка cache: true/false для ассетов - Логотипы команд: cache: false (часто меняются) - Логотип канала: cache: true (редко меняется) - Видео-пак: cache: true (редко меняется) - Версия 0.2.6
This commit is contained in:
179
src/processor.rs
179
src/processor.rs
@@ -295,35 +295,31 @@ async fn generate_and_submit_jobs(
|
||||
) -> Result<Vec<(String, String)>> {
|
||||
info!("Preparing Nexrender jobs...");
|
||||
|
||||
// ========== ЗАГРУЖАЕМ SPORT PACKS ==========
|
||||
let sport_sheet = workbook
|
||||
.get_sheet("SPORT")
|
||||
.ok_or_else(|| anyhow!("Sheet 'SPORT' not found"))?;
|
||||
// ========== ЗАГРУЖАЕМ SPORT PACKS ==========
|
||||
let sport_sheet = workbook
|
||||
.get_sheet("SPORT")
|
||||
.ok_or_else(|| anyhow!("Sheet 'SPORT' not found"))?;
|
||||
|
||||
let packs: HashMap<String, String> = sport_sheet
|
||||
.rows
|
||||
.iter()
|
||||
.filter_map(|row| Some((row.get("SPORT")?.clone(), row.get("LINK")?.clone())))
|
||||
.collect();
|
||||
info!("Loaded {} sport packs", packs.len());
|
||||
// sport_sheet автоматически освободится при выходе из функции
|
||||
|
||||
// Лист SPORT больше не нужен
|
||||
let sport_rows_count = sport_sheet.rows.len();
|
||||
debug!("Dropped SPORT sheet ({} rows)", sport_rows_count);
|
||||
let packs: HashMap<String, String> = sport_sheet
|
||||
.rows
|
||||
.iter()
|
||||
.filter_map(|row| Some((row.get("SPORT")?.clone(), row.get("LINK")?.clone())))
|
||||
.collect();
|
||||
info!("Loaded {} sport packs", packs.len());
|
||||
debug!("Processed SPORT sheet ({} rows)", sport_sheet.rows.len());
|
||||
|
||||
// ========== ЗАГРУЖАЕМ TEAM LOGOS ==========
|
||||
let teams_sheet = workbook
|
||||
.get_sheet("TEAMS")
|
||||
.ok_or_else(|| anyhow!("Sheet 'TEAMS' not found"))?;
|
||||
|
||||
|
||||
debug!("TEAMS headers: {:?}", teams_sheet.headers);
|
||||
if log_enabled!(Level::Debug) {
|
||||
for (i, row) in teams_sheet.rows.iter().take(3).enumerate() {
|
||||
debug!("TEAMS row {}: {:?}", i, row);
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
let mut logos = LogoRegistry::with_capacity(teams_sheet.rows.len());
|
||||
for row in &teams_sheet.rows {
|
||||
let team = get_cell_fuzzy(row, "TEAM");
|
||||
@@ -336,7 +332,7 @@ info!("Loaded {} sport packs", packs.len());
|
||||
}
|
||||
let total_teams = teams_sheet.rows.len();
|
||||
info!("Loaded {} team logos", total_teams);
|
||||
|
||||
|
||||
// Дебаг: выводим статистику по TEAMS
|
||||
if log_enabled!(Level::Debug) {
|
||||
let sports: HashSet<_> = teams_sheet
|
||||
@@ -357,32 +353,30 @@ info!("Loaded {} sport packs", packs.len());
|
||||
debug!(" TEAMS examples for '{}': {:?}", sport, examples);
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
// ========== ЗАГРУЖАЕМ CHANNEL LOGOS ==========
|
||||
let channel_sheet = workbook
|
||||
.get_sheet("CHANELL")
|
||||
.ok_or_else(|| anyhow!("Sheet 'CHANELL' not found"))?;
|
||||
|
||||
|
||||
let channels: HashMap<String, String> = channel_sheet
|
||||
.rows
|
||||
.iter()
|
||||
.filter_map(|row| Some((row.get("CHANELL")?.clone(), row.get("LINK")?.clone())))
|
||||
.collect();
|
||||
info!("Loaded {} channel logos", channels.len());
|
||||
|
||||
|
||||
if log_enabled!(Level::Debug) {
|
||||
debug!("Channels loaded: {:?}", channels.keys().collect::<Vec<_>>());
|
||||
}
|
||||
|
||||
|
||||
// ========== ОБРАБАТЫВАЕМ START (только активные строки) ==========
|
||||
let start_sheet = workbook
|
||||
.get_sheet("Start")
|
||||
.ok_or_else(|| anyhow!("Sheet 'Start' not found"))?;
|
||||
|
||||
|
||||
let total_start_rows = start_sheet.rows.len();
|
||||
|
||||
|
||||
// Сразу фильтруем только строки с STATE = "FALSE"
|
||||
let active_rows: Vec<(usize, &HashMap<String, String>)> = start_sheet
|
||||
.rows
|
||||
@@ -390,26 +384,26 @@ info!("Loaded {} sport packs", packs.len());
|
||||
.enumerate()
|
||||
.filter(|(_, row)| row.get("STATE").map(|s| s.as_str()) == Some("FALSE"))
|
||||
.collect();
|
||||
|
||||
|
||||
info!(
|
||||
"Found {} active rows (STATE='FALSE') out of {} total",
|
||||
active_rows.len(),
|
||||
"Found {} active rows (STATE='FALSE') out of {} total",
|
||||
active_rows.len(),
|
||||
total_start_rows
|
||||
);
|
||||
|
||||
|
||||
// Собираем использованные команды для очистки logos
|
||||
let mut used_teams: HashSet<(String, String)> = HashSet::new();
|
||||
let mut jobs: Vec<JobData> = Vec::with_capacity(active_rows.len() * 3);
|
||||
|
||||
|
||||
// Для дебага
|
||||
let mut missing_teams: HashSet<String> = HashSet::new();
|
||||
let mut missing_sports: HashSet<String> = HashSet::new();
|
||||
|
||||
|
||||
for (idx, row) in active_rows {
|
||||
let team_a = row.get("TEAM A").cloned().unwrap_or_default();
|
||||
let team_b = row.get("TEAM B").cloned().unwrap_or_default();
|
||||
let sport = row.get("SPORT").cloned().unwrap_or_default();
|
||||
|
||||
|
||||
// Запоминаем использованные команды
|
||||
if !team_a.is_empty() {
|
||||
used_teams.insert((team_a.clone(), sport.clone()));
|
||||
@@ -417,7 +411,7 @@ info!("Loaded {} sport packs", packs.len());
|
||||
if !team_b.is_empty() {
|
||||
used_teams.insert((team_b.clone(), sport.clone()));
|
||||
}
|
||||
|
||||
|
||||
if log_enabled!(Level::Debug) {
|
||||
let logo_a = logos.find(&team_a, &sport);
|
||||
let logo_b = logos.find(&team_b, &sport);
|
||||
@@ -431,21 +425,20 @@ info!("Loaded {} sport packs", packs.len());
|
||||
missing_sports.insert(sport.clone());
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
if let Some(job) = JobData::from_row(row, idx, &packs, &logos, &channels) {
|
||||
jobs.extend(job.create_variants());
|
||||
}
|
||||
}
|
||||
|
||||
|
||||
// Очищаем logos от неиспользуемых команд
|
||||
let logos_before = logos.len();
|
||||
logos.retain(|team, sport| used_teams.contains(&(team.to_string(), sport.to_string())));
|
||||
info!(
|
||||
"Retained {} used team logos (cleaned up {} unused)",
|
||||
"Retained {} used team logos (cleaned up {} unused)",
|
||||
logos.len(),
|
||||
logos_before - logos.len()
|
||||
);
|
||||
|
||||
|
||||
// Дебаг: выводим сводку по отсутствующим логотипам
|
||||
if log_enabled!(Level::Debug) && !missing_teams.is_empty() {
|
||||
@@ -461,75 +454,71 @@ info!("Loaded {} sport packs", packs.len());
|
||||
|
||||
info!("Generated {} total jobs (including variants)", jobs.len());
|
||||
|
||||
// ========== СОРТИРОВКА ПО ДАТЕ И ВРЕМЕНИ ==========
|
||||
jobs.sort_by(|a, b| match a.sort_date.cmp(&b.sort_date) {
|
||||
std::cmp::Ordering::Equal => match a.sort_time.cmp(&b.sort_time) {
|
||||
std::cmp::Ordering::Equal => a.variant_order.cmp(&b.variant_order),
|
||||
other => other,
|
||||
},
|
||||
other => other,
|
||||
});
|
||||
|
||||
info!("Jobs sorted by date and time (earliest first)");
|
||||
|
||||
if log_enabled!(Level::Debug) {
|
||||
for (i, job) in jobs.iter().take(10).enumerate() {
|
||||
debug!(
|
||||
" {}: {} {} - {} (variant: {})",
|
||||
i + 1,
|
||||
job.sort_date,
|
||||
job.sort_time,
|
||||
job.outfile_name,
|
||||
job.variant_order
|
||||
);
|
||||
}
|
||||
}
|
||||
|
||||
if jobs.is_empty() {
|
||||
info!("No jobs with STATE='FALSE' found");
|
||||
return Ok(Vec::new());
|
||||
}
|
||||
|
||||
// ========== ОТПРАВКА ЗАДАНИЙ ==========
|
||||
cleanup_finished_jobs(&config.nexrender_api_url).await?;
|
||||
// ========== ОТПРАВКА ЗАДАНИЙ (ПОСЛЕДОВАТЕЛЬНО) ==========
|
||||
cleanup_finished_jobs(&config.nexrender_api_url).await?;
|
||||
|
||||
let http_client = Client::new();
|
||||
let mut submitted_jobs: Vec<(String, String)> = Vec::with_capacity(jobs.len());
|
||||
let http_client = Client::new();
|
||||
let mut submitted_jobs: Vec<(String, String)> = Vec::with_capacity(jobs.len());
|
||||
|
||||
// Запускаем ВСЕ группы параллельно
|
||||
let mut all_tasks = Vec::new();
|
||||
for job in jobs {
|
||||
let nexrender_job = job.to_nexrender_job(config);
|
||||
info!("Submitting job: {}", job.outfile_name);
|
||||
|
||||
for chunk in jobs.chunks(3) {
|
||||
for job in chunk {
|
||||
let client = http_client.clone();
|
||||
let api_url = config.nexrender_api_url.clone();
|
||||
let config_clone = config.clone();
|
||||
let job_owned = job.clone();
|
||||
|
||||
all_tasks.push(tokio::spawn(async move {
|
||||
let nexrender_job = job_owned.to_nexrender_job(&config_clone);
|
||||
info!("Submitting job: {}", job_owned.outfile_name);
|
||||
|
||||
if log_enabled!(Level::Debug) {
|
||||
debug!(
|
||||
"Job details: sport='{}', league='{}', channel='{}', team_a='{}', team_b='{}'",
|
||||
job_owned.sport, job_owned.league, job_owned.channel, job_owned.team_a, job_owned.team_b
|
||||
);
|
||||
if log_enabled!(Level::Debug) {
|
||||
debug!(
|
||||
"Job details: sport='{}', league='{}', channel='{}', team_a='{}', team_b='{}'",
|
||||
job.sport, job.league, job.channel, job.team_a, job.team_b
|
||||
);
|
||||
}
|
||||
|
||||
let response = http_client
|
||||
.post(&config.nexrender_api_url)
|
||||
.json(&nexrender_job)
|
||||
.send()
|
||||
.await?;
|
||||
|
||||
if response.status().is_success() {
|
||||
let result: Value = response.json().await?;
|
||||
if let Some(uid) = result.get("uid").and_then(|u| u.as_str()) {
|
||||
info!("Job submitted: {} (UID: {})", job.outfile_name, uid);
|
||||
submitted_jobs.push((uid.to_string(), job.outfile_name.clone()));
|
||||
}
|
||||
|
||||
let response = client.post(&api_url).json(&nexrender_job).send().await?;
|
||||
|
||||
if response.status().is_success() {
|
||||
let result: Value = response.json().await?;
|
||||
if let Some(uid) = result.get("uid").and_then(|u| u.as_str()) {
|
||||
Ok::<_, anyhow::Error>(Some((uid.to_string(), job_owned.outfile_name)))
|
||||
} else {
|
||||
Ok(None)
|
||||
}
|
||||
} else {
|
||||
let status = response.status();
|
||||
let text = response.text().await.unwrap_or_default();
|
||||
error!("Failed to submit job ({}): {}", status, text);
|
||||
Ok(None)
|
||||
}
|
||||
}));
|
||||
}
|
||||
}
|
||||
|
||||
// Ждём завершения ВСЕХ задач
|
||||
for task in all_tasks {
|
||||
match task.await {
|
||||
Ok(Ok(Some((uid, outfile_name)))) => {
|
||||
info!("Job submitted: {} (UID: {})", outfile_name, uid);
|
||||
submitted_jobs.push((uid, outfile_name));
|
||||
}
|
||||
Ok(Ok(None)) => {
|
||||
debug!("Job submission returned no UID");
|
||||
}
|
||||
Ok(Err(e)) => {
|
||||
error!("Job submission error: {}", e);
|
||||
}
|
||||
Err(e) => {
|
||||
error!("Task join error: {}", e);
|
||||
} else {
|
||||
let status = response.status();
|
||||
let text = response.text().await.unwrap_or_default();
|
||||
error!("Failed to submit job ({}): {}", status, text);
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
info!(
|
||||
"Successfully submitted {} jobs to Nexrender",
|
||||
submitted_jobs.len()
|
||||
@@ -568,4 +557,4 @@ pub async fn fetch_all_jobs(api_url: &str) -> Result<Vec<Value>> {
|
||||
} else {
|
||||
Ok(Vec::new())
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user