feat: групповая отправка заданий по 3 для сохранения порядка в очереди

- Отправка заданий группами по 3 (оригинал + today + tomorrow) параллельно
- Сохранение порядка внутри группы для корректной FIFO очереди Nexrender
- Добавлен UID в логи отправки заданий
- Версия обновлена до 0.2.4
This commit is contained in:
2026-04-17 18:29:04 +03:00
parent 6f13679dd9
commit 9929dab66f
4 changed files with 66 additions and 63 deletions

View File

@@ -303,13 +303,11 @@ async fn generate_and_submit_jobs(
// Вспомогательная функция для поиска значения по части ключа
fn get_cell_fuzzy(row: &HashMap<String, String>, key_part: &str) -> Option<String> {
// Точное совпадение
if let Some(val) = row.get(key_part) {
if !val.is_empty() {
return Some(val.clone());
}
}
// Поиск по содержанию (без учёта регистра)
for (k, v) in row.iter() {
if k.to_lowercase().contains(&key_part.to_lowercase()) && !v.is_empty() {
return Some(v.clone());
@@ -427,55 +425,61 @@ async fn generate_and_submit_jobs(
let http_client = Client::new();
let mut submitted_jobs: Vec<(String, String)> = Vec::with_capacity(jobs.len());
let mut tasks = Vec::with_capacity(jobs.len());
for job in jobs {
let client = http_client.clone();
let api_url = config.nexrender_api_url.clone();
let config_clone = config.clone();
tasks.push(tokio::spawn(async move {
let nexrender_job = job.to_nexrender_job(&config_clone);
info!("Submitting job: {}", job.outfile_name);
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 = 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.outfile_name)))
// Отправляем задания группами по 3 (оригинал + today + tomorrow)
for chunk in jobs.chunks(3) {
let mut tasks = Vec::with_capacity(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();
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
);
}
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)
}
} 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 tasks {
match task.await {
Ok(Ok(Some((uid, outfile_name)))) => {
info!("Job submitted: {}", outfile_name);
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);
}));
}
// Ждём завершения группы
for task in 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);
}
}
}
}
@@ -518,4 +522,4 @@ pub async fn fetch_all_jobs(api_url: &str) -> Result<Vec<Value>> {
} else {
Ok(Vec::new())
}
}
}