This commit is contained in:
2026-04-17 11:58:16 +03:00
parent 4f1c6d1ab7
commit 8721353b8d
6 changed files with 1093 additions and 463 deletions

View File

@@ -1,480 +1,42 @@
mod config;
mod nexrender;
mod synology;
mod processor;
mod web;
use anyhow::{anyhow, Result};
use calamine::{Data, Reader, Xlsx};
use chrono::{Duration, NaiveDate};
use anyhow::Result;
use clap::Parser;
use config::Config;
use log::{debug, error, info};
use nexrender::{JobData, LogoRegistry};
use reqwest::Client;
use serde_json::Value;
use std::collections::HashMap;
use std::io::Cursor;
use std::path::Path;
use std::time::Duration as StdDuration;
use synology::SynologyClient;
use tokio::time::sleep;
use log::info;
// Структуры для in-memory данных
#[derive(Debug, Clone)]
pub struct SheetData {
pub name: String,
pub headers: Vec<String>,
pub rows: Vec<HashMap<String, String>>,
}
#[derive(Parser)]
#[command(author, version, about, long_about = None)]
struct Cli {
/// Run web server instead of one-time processing
#[arg(short, long)]
web: bool,
#[derive(Debug, Default)]
pub struct ExcelWorkbook {
pub sheets: Vec<SheetData>,
/// Run one-time processing (default if no flags)
#[arg(short, long)]
once: bool,
}
#[tokio::main]
async fn main() -> Result<()> {
env_logger::Builder::from_env(env_logger::Env::default().default_filter_or("info")).init();
info!("Starting AE Anons processor v0.1.1 (in-memory mode)");
let cli = Cli::parse();
let config = Config::from_env()?;
debug!("Configuration loaded");
let mut client = SynologyClient::new(&config.nas_fqdn);
client.login(&config.nas_user, &config.nas_pass).await?;
info!("Successfully authenticated with Synology NAS");
let info = client.get_info().await?;
info!("Connected to NAS: {}", info.hostname);
// Скачиваем и парсим Excel в памяти
let workbook = download_and_parse_excel_in_memory(&mut client, &config).await?;
display_workbook_structure(&workbook);
process_nexrender_jobs(&workbook, &config).await?;
client.logout().await?;
info!("Session terminated successfully");
Ok(())
}
async fn download_and_parse_excel_in_memory(
client: &mut SynologyClient,
config: &Config,
) -> Result<ExcelWorkbook> {
let file_name_full = Path::new(&config.nas_file)
.file_name()
.and_then(|n| n.to_str())
.unwrap_or("unknown");
let search_name = file_name_full.replace(".osheet", "");
let expected_path = Path::new(&config.nas_file)
.parent()
.and_then(|p| p.to_str())
.unwrap_or("");
let api_format_path = expected_path
.replace("Team Folder", "team-folders")
.replace(' ', "");
let full_expected_path_api = format!("{}/{}", api_format_path, file_name_full);
info!("Searching for file: {}", full_expected_path_api);
let search_result = client.search_file_by_name(&search_name).await?;
let items = search_result["data"]["items"]
.as_array()
.ok_or_else(|| anyhow!("No items in search result"))?;
if items.is_empty() {
return Err(anyhow!("File '{}' not found", file_name_full));
}
let exact_match: Vec<&Value> = items
.iter()
.filter(|item| {
item.get("display_path")
.and_then(|p| p.as_str())
.unwrap_or("")
== full_expected_path_api
})
.collect();
if exact_match.is_empty() {
return Err(anyhow!("File not found at: {}", full_expected_path_api));
}
let file = exact_match[0];
let file_id = file["file_id"]
.as_str()
.ok_or_else(|| anyhow!("Missing file_id"))?;
let actual_file_name = file["name"]
.as_str()
.ok_or_else(|| anyhow!("Missing name"))?;
info!("Found file: {}", actual_file_name);
info!("Exporting file from Synology Office to Excel format (in-memory)...");
// Получаем бинарные данные Excel напрямую в память
let excel_data = client.export_by_file_id(file_id, actual_file_name).await?;
info!("Exported {} bytes to memory", excel_data.len());
info!("Parsing Excel workbook from memory...");
parse_excel_from_bytes(&excel_data)
}
fn parse_excel_from_bytes(data: &[u8]) -> Result<ExcelWorkbook> {
let cursor = Cursor::new(data);
let mut workbook: Xlsx<_> = calamine::open_workbook_from_rs(cursor)
.map_err(|e| anyhow!("Failed to open workbook from memory: {}", e))?;
let sheet_names = workbook.sheet_names().to_vec();
let mut excel_workbook = ExcelWorkbook::default();
for sheet_name in sheet_names {
debug!("Processing sheet: '{}'", sheet_name);
let range = workbook
.worksheet_range(&sheet_name)
.map_err(|e| anyhow!("Failed to read sheet '{}': {}", sheet_name, e))?;
let sheet_data = parse_sheet_dynamic_optimized(&sheet_name, &range)?;
excel_workbook.sheets.push(sheet_data);
}
Ok(excel_workbook)
}
fn parse_sheet_dynamic_optimized(
sheet_name: &str,
range: &calamine::Range<Data>,
) -> Result<SheetData> {
// Предварительно выделяем память для избежания реаллокаций
let (row_count, col_count) = range.get_size();
let mut data_matrix: Vec<Vec<String>> = Vec::with_capacity(row_count);
// Строим матрицу данных (без enumerate)
for row in range.rows() {
let mut row_data = Vec::with_capacity(col_count);
for cell in row {
row_data.push(cell_to_string_optimized(cell));
}
data_matrix.push(row_data);
}
if data_matrix.is_empty() {
return Ok(SheetData {
name: sheet_name.to_string(),
headers: Vec::new(),
rows: Vec::new(),
});
}
// Извлекаем заголовки с дедупликацией
let headers: Vec<String> = data_matrix[0]
.iter()
.enumerate()
.map(|(idx, header)| {
let h = header.trim().to_string();
if h.is_empty() {
format!("Column_{}", idx + 1)
} else {
h
}
})
.collect();
// Парсим строки данных
let rows_data: Vec<HashMap<String, String>> = data_matrix
.iter()
.skip(1)
.filter_map(|row_values| {
let mut row_map = HashMap::with_capacity(headers.len());
let mut has_data = false;
for (col_idx, value) in row_values.iter().enumerate() {
if col_idx < headers.len() && !value.is_empty() {
row_map.insert(headers[col_idx].clone(), value.clone());
has_data = true;
}
}
if has_data {
Some(row_map)
} else {
None
}
})
.take(10000) // Ограничение для безопасности
.collect();
Ok(SheetData {
name: sheet_name.to_string(),
headers,
rows: rows_data,
})
}
#[inline]
fn cell_to_string_optimized(cell: &Data) -> String {
match cell {
Data::Empty => String::new(),
Data::String(s) => s.clone(),
Data::Float(f) => {
if is_excel_date(*f) {
excel_date_to_string(*f)
} else if f.fract() == 0.0 {
// Используем itoa для целых чисел (опционально)
format!("{:.0}", f)
} else {
// Используем ryu для float (опционально)
f.to_string()
}
}
Data::Int(i) => {
let f = *i as f64;
if is_excel_date(f) {
excel_date_to_string(f)
} else {
i.to_string()
}
}
Data::Bool(b) => {
if *b {
"TRUE".to_string()
} else {
"FALSE".to_string()
}
}
Data::DateTime(dt) => {
let serial = dt.as_f64();
if is_excel_date(serial) {
excel_date_to_string(serial)
} else {
dt.to_string()
}
}
Data::DateTimeIso(s) => s.clone(),
Data::DurationIso(s) => s.clone(),
Data::Error(e) => format!("{:?}", e),
}
}
#[inline]
fn is_excel_date(value: f64) -> bool {
(1.0..100000.0).contains(&value)
}
fn excel_date_to_string(serial: f64) -> String {
let days = serial as i64;
let base = NaiveDate::from_ymd_opt(1899, 12, 30).unwrap();
if let Some(date) = base.checked_add_signed(Duration::days(days)) {
// Используем метод format напрямую - он публичный
date.format("%d.%m.%Y").to_string()
if cli.web {
info!("Starting AE Anons web server v{}", env!("CARGO_PKG_VERSION"));
web::run_web_server(config).await?;
} else {
serial.to_string()
}
}
fn display_workbook_structure(workbook: &ExcelWorkbook) {
info!("Workbook contains {} sheets", workbook.sheets.len());
for (idx, sheet) in workbook.sheets.iter().enumerate() {
debug!(
"Sheet #{}: '{}' | Headers: {} | Rows: {}",
idx + 1,
sheet.name,
sheet.headers.len(),
sheet.rows.len()
);
}
}
async fn process_nexrender_jobs(workbook: &ExcelWorkbook, config: &Config) -> Result<()> {
info!("Preparing Nexrender jobs...");
let start_sheet = workbook
.get_sheet("Start")
.ok_or_else(|| anyhow!("Sheet 'Start' not found"))?;
let sport_sheet = workbook
.get_sheet("SPORT")
.ok_or_else(|| anyhow!("Sheet 'SPORT' not found"))?;
let teams_sheet = workbook
.get_sheet("TEAMS")
.ok_or_else(|| anyhow!("Sheet 'TEAMS' not found"))?;
let channel_sheet = workbook
.get_sheet("CHANELL")
.ok_or_else(|| anyhow!("Sheet 'CHANELL' 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());
let mut logos = LogoRegistry::with_capacity(teams_sheet.rows.len());
for row in &teams_sheet.rows {
if let (Some(team), Some(sport), Some(link)) =
(row.get("TEAM"), row.get("SPORT"), row.get("LINK"))
{
logos.insert(team.clone(), sport.clone(), link.clone());
}
}
info!("Loaded {} team logos", teams_sheet.rows.len());
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());
// Предварительно выделяем память для jobs
let mut jobs: Vec<JobData> = Vec::with_capacity(start_sheet.rows.len() * 3);
for (idx, row) in start_sheet.rows.iter().enumerate() {
if let Some(state) = row.get("STATE") {
if state == "FALSE" {
if let Some(job) = JobData::from_row(row, idx, &packs, &logos, &channels) {
jobs.extend(job.create_variants());
}
}
}
}
info!("Generated {} total jobs (including variants)", jobs.len());
if jobs.is_empty() {
info!("No jobs with STATE='FALSE' found");
return Ok(());
}
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());
// Отправляем jobs параллельно для ускорения
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); // ✅ Передаём Config
info!("Submitting job: {}", job.outfile_name);
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)))
} 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 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);
}
}
}
info!(
"Successfully submitted {} jobs to Nexrender",
submitted_jobs.len()
);
if !submitted_jobs.is_empty() {
info!("Monitoring job completion...");
monitor_jobs(&http_client, &config.nexrender_api_url, submitted_jobs).await?;
// По умолчанию или с флагом --once выполняем однократную обработку
info!("Starting AE Anons processor v{} (one-time mode)", env!("CARGO_PKG_VERSION"));
let submitted = processor::process_spreadsheet(&config).await?;
info!("Submitted {} jobs. Exiting.", submitted.len());
}
Ok(())
}
// ... остальные функции без изменений ...
impl ExcelWorkbook {
pub fn get_sheet(&self, name: &str) -> Option<&SheetData> {
self.sheets.iter().find(|s| s.name == name)
}
}
async fn cleanup_finished_jobs(api_url: &str) -> Result<()> {
let client = Client::new();
let response = client.get(api_url).send().await?;
if response.status().is_success() {
let jobs: Vec<Value> = response.json().await?;
for job in jobs {
if let (Some(uid), Some(state)) = (
job.get("uid").and_then(|u| u.as_str()),
job.get("state").and_then(|s| s.as_str()),
) {
if state == "finished" || state == "error" {
let _ = client.delete(&format!("{}/{}", api_url, uid)).send().await;
info!("Cleaned up completed job: {}", uid);
}
}
}
}
Ok(())
}
async fn monitor_jobs(
client: &Client,
api_url: &str,
mut pending: Vec<(String, String)>,
) -> Result<()> {
while !pending.is_empty() {
sleep(StdDuration::from_secs(25)).await;
let mut remaining = Vec::with_capacity(pending.len());
for (uid, outname) in pending.drain(..) {
let response = client.get(&format!("{}/{}", api_url, uid)).send().await?;
if response.status().is_success() {
let job: Value = response.json().await?;
let state = job
.get("state")
.and_then(|s| s.as_str())
.unwrap_or("unknown");
match state {
"finished" => {
info!("Job completed: {}", outname)
}
"error" => error!("Job failed: {}", outname),
_ => remaining.push((uid, outname)),
}
} else {
remaining.push((uid, outname));
}
}
pending = remaining;
if !pending.is_empty() {
info!("Waiting for {} jobs to complete...", pending.len());
}
}
info!("All jobs completed successfully");
Ok(())
}
}

434
src/processor.rs Normal file
View File

@@ -0,0 +1,434 @@
use crate::config::Config;
use crate::nexrender::{JobData, LogoRegistry};
use crate::synology::SynologyClient;
use anyhow::{anyhow, Result};
use calamine::{Data, Reader, Xlsx};
use chrono::{Duration, NaiveDate};
use log::{debug, error, info};
use reqwest::Client;
use serde_json::Value;
use std::collections::HashMap;
use std::io::Cursor;
use std::path::Path;
// Добавить после импортов в processor.rs
#[derive(Debug, Clone)]
pub struct SheetData {
pub name: String,
pub headers: Vec<String>,
pub rows: Vec<HashMap<String, String>>,
}
#[derive(Debug, Default)]
pub struct ExcelWorkbook {
pub sheets: Vec<SheetData>,
}
impl ExcelWorkbook {
pub fn get_sheet(&self, name: &str) -> Option<&SheetData> {
self.sheets.iter().find(|s| s.name == name)
}
}
/// Основная функция обработки: скачивает Excel, генерирует задания, отправляет в Nexrender
pub async fn process_spreadsheet(config: &Config) -> Result<Vec<(String, String)>> {
let mut client = SynologyClient::new(&config.nas_fqdn);
client.login(&config.nas_user, &config.nas_pass).await?;
info!("Successfully authenticated with Synology NAS");
let info = client.get_info().await?;
info!("Connected to NAS: {}", info.hostname);
// Скачиваем и парсим Excel в памяти
let workbook = download_and_parse_excel_in_memory(&mut client, config).await?;
display_workbook_structure(&workbook);
let submitted = generate_and_submit_jobs(&workbook, config).await?;
client.logout().await?;
info!("Session terminated successfully");
Ok(submitted)
}
async fn download_and_parse_excel_in_memory(
client: &mut SynologyClient,
config: &Config,
) -> Result<ExcelWorkbook> {
let file_name_full = Path::new(&config.nas_file)
.file_name()
.and_then(|n| n.to_str())
.unwrap_or("unknown");
let search_name = file_name_full.replace(".osheet", "");
let expected_path = Path::new(&config.nas_file)
.parent()
.and_then(|p| p.to_str())
.unwrap_or("");
let api_format_path = expected_path
.replace("Team Folder", "team-folders")
.replace(' ', "");
let full_expected_path_api = format!("{}/{}", api_format_path, file_name_full);
info!("Searching for file: {}", full_expected_path_api);
let search_result = client.search_file_by_name(&search_name).await?;
let items = search_result["data"]["items"]
.as_array()
.ok_or_else(|| anyhow!("No items in search result"))?;
if items.is_empty() {
return Err(anyhow!("File '{}' not found", file_name_full));
}
let exact_match: Vec<&Value> = items
.iter()
.filter(|item| {
item.get("display_path")
.and_then(|p| p.as_str())
.unwrap_or("")
== full_expected_path_api
})
.collect();
if exact_match.is_empty() {
return Err(anyhow!("File not found at: {}", full_expected_path_api));
}
let file = exact_match[0];
let file_id = file["file_id"]
.as_str()
.ok_or_else(|| anyhow!("Missing file_id"))?;
let actual_file_name = file["name"]
.as_str()
.ok_or_else(|| anyhow!("Missing name"))?;
info!("Found file: {}", actual_file_name);
info!("Exporting file from Synology Office to Excel format (in-memory)...");
// Получаем бинарные данные Excel напрямую в память
let excel_data = client.export_by_file_id(file_id, actual_file_name).await?;
info!("Exported {} bytes to memory", excel_data.len());
info!("Parsing Excel workbook from memory...");
parse_excel_from_bytes(&excel_data)
}
fn parse_excel_from_bytes(data: &[u8]) -> Result<ExcelWorkbook> {
let cursor = Cursor::new(data);
let mut workbook: Xlsx<_> = calamine::open_workbook_from_rs(cursor)
.map_err(|e| anyhow!("Failed to open workbook from memory: {}", e))?;
let sheet_names = workbook.sheet_names().to_vec();
let mut excel_workbook = ExcelWorkbook::default();
for sheet_name in sheet_names {
debug!("Processing sheet: '{}'", sheet_name);
let range = workbook
.worksheet_range(&sheet_name)
.map_err(|e| anyhow!("Failed to read sheet '{}': {}", sheet_name, e))?;
let sheet_data = parse_sheet_dynamic_optimized(&sheet_name, &range)?;
excel_workbook.sheets.push(sheet_data);
}
Ok(excel_workbook)
}
fn parse_sheet_dynamic_optimized(
sheet_name: &str,
range: &calamine::Range<Data>,
) -> Result<SheetData> {
let (row_count, col_count) = range.get_size();
let mut data_matrix: Vec<Vec<String>> = Vec::with_capacity(row_count);
for row in range.rows() {
let mut row_data = Vec::with_capacity(col_count);
for cell in row {
row_data.push(cell_to_string_optimized(cell));
}
data_matrix.push(row_data);
}
if data_matrix.is_empty() {
return Ok(SheetData {
name: sheet_name.to_string(),
headers: Vec::new(),
rows: Vec::new(),
});
}
// Извлекаем заголовки с дедупликацией
let headers: Vec<String> = data_matrix[0]
.iter()
.enumerate()
.map(|(idx, header)| {
let h = header.trim().to_string();
if h.is_empty() {
format!("Column_{}", idx + 1)
} else {
h
}
})
.collect();
// Парсим строки данных
let rows_data: Vec<HashMap<String, String>> = data_matrix
.iter()
.skip(1)
.filter_map(|row_values| {
let mut row_map = HashMap::with_capacity(headers.len());
let mut has_data = false;
for (col_idx, value) in row_values.iter().enumerate() {
if col_idx < headers.len() && !value.is_empty() {
row_map.insert(headers[col_idx].clone(), value.clone());
has_data = true;
}
}
if has_data {
Some(row_map)
} else {
None
}
})
.take(10000) // Ограничение для безопасности
.collect();
Ok(SheetData {
name: sheet_name.to_string(),
headers,
rows: rows_data,
})
}
#[inline]
fn cell_to_string_optimized(cell: &Data) -> String {
match cell {
Data::Empty => String::new(),
Data::String(s) => s.clone(),
Data::Float(f) => {
if is_excel_date(*f) {
excel_date_to_string(*f)
} else if f.fract() == 0.0 {
// Используем itoa для целых чисел (опционально)
format!("{:.0}", f)
} else {
// Используем ryu для float (опционально)
f.to_string()
}
}
Data::Int(i) => {
let f = *i as f64;
if is_excel_date(f) {
excel_date_to_string(f)
} else {
i.to_string()
}
}
Data::Bool(b) => {
if *b {
"TRUE".to_string()
} else {
"FALSE".to_string()
}
}
Data::DateTime(dt) => {
let serial = dt.as_f64();
if is_excel_date(serial) {
excel_date_to_string(serial)
} else {
dt.to_string()
}
}
Data::DateTimeIso(s) => s.clone(),
Data::DurationIso(s) => s.clone(),
Data::Error(e) => format!("{:?}", e),
}
}
#[inline]
fn is_excel_date(value: f64) -> bool {
(1.0..100000.0).contains(&value)
}
fn excel_date_to_string(serial: f64) -> String {
let days = serial as i64;
let base = NaiveDate::from_ymd_opt(1899, 12, 30).unwrap();
if let Some(date) = base.checked_add_signed(Duration::days(days)) {
// Используем метод format напрямую - он публичный
date.format("%d.%m.%Y").to_string()
} else {
serial.to_string()
}
}
fn display_workbook_structure(workbook: &ExcelWorkbook) {
info!("Workbook contains {} sheets", workbook.sheets.len());
for (idx, sheet) in workbook.sheets.iter().enumerate() {
debug!(
"Sheet #{}: '{}' | Headers: {} | Rows: {}",
idx + 1,
sheet.name,
sheet.headers.len(),
sheet.rows.len()
);
}
}
async fn generate_and_submit_jobs(
workbook: &ExcelWorkbook,
config: &Config,
) -> Result<Vec<(String, String)>> {
info!("Preparing Nexrender jobs...");
let start_sheet = workbook
.get_sheet("Start")
.ok_or_else(|| anyhow!("Sheet 'Start' not found"))?;
let sport_sheet = workbook
.get_sheet("SPORT")
.ok_or_else(|| anyhow!("Sheet 'SPORT' not found"))?;
let teams_sheet = workbook
.get_sheet("TEAMS")
.ok_or_else(|| anyhow!("Sheet 'TEAMS' not found"))?;
let channel_sheet = workbook
.get_sheet("CHANELL")
.ok_or_else(|| anyhow!("Sheet 'CHANELL' 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());
let mut logos = LogoRegistry::with_capacity(teams_sheet.rows.len());
for row in &teams_sheet.rows {
if let (Some(team), Some(sport), Some(link)) =
(row.get("TEAM"), row.get("SPORT"), row.get("LINK"))
{
logos.insert(team.clone(), sport.clone(), link.clone());
}
}
info!("Loaded {} team logos", teams_sheet.rows.len());
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());
let mut jobs: Vec<JobData> = Vec::with_capacity(start_sheet.rows.len() * 3);
for (idx, row) in start_sheet.rows.iter().enumerate() {
if let Some(state) = row.get("STATE") {
if state == "FALSE" {
if let Some(job) = JobData::from_row(row, idx, &packs, &logos, &channels) {
jobs.extend(job.create_variants());
}
}
}
}
info!("Generated {} total jobs (including variants)", jobs.len());
if jobs.is_empty() {
info!("No jobs with STATE='FALSE' found");
return Ok(Vec::new());
}
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 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);
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)))
} 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 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);
}
}
}
info!(
"Successfully submitted {} jobs to Nexrender",
submitted_jobs.len()
);
Ok(submitted_jobs)
}
pub async fn cleanup_finished_jobs(api_url: &str) -> Result<()> {
let client = Client::new();
let response = client.get(api_url).send().await?;
if response.status().is_success() {
let jobs: Vec<Value> = response.json().await?;
for job in jobs {
if let (Some(uid), Some(state)) = (
job.get("uid").and_then(|u| u.as_str()),
job.get("state").and_then(|s| s.as_str()),
) {
if state == "finished" || state == "error" {
let _ = client.delete(&format!("{}/{}", api_url, uid)).send().await;
info!("Cleaned up completed job: {}", uid);
}
}
}
}
Ok(())
}
pub async fn fetch_all_jobs(api_url: &str) -> Result<Vec<Value>> {
let client = Client::new();
let response = client.get(api_url).send().await?;
if response.status().is_success() {
let jobs: Vec<Value> = response.json().await?;
Ok(jobs)
} else {
Ok(Vec::new())
}
}

214
src/static/index.html Normal file
View File

@@ -0,0 +1,214 @@
<!DOCTYPE html>
<html>
<head>
<meta charset="utf-8">
<title>AE Anons Control Panel</title>
<style>
body { font-family: -apple-system, BlinkMacSystemFont, 'Segoe UI', Roboto, sans-serif; margin: 20px; background: #f5f5f5; }
.container { max-width: 1400px; margin: 0 auto; background: white; padding: 20px; border-radius: 8px; box-shadow: 0 2px 4px rgba(0,0,0,0.1); }
h1 { color: #333; margin-top: 0; }
table { border-collapse: collapse; width: 100%; margin-top: 20px; }
th, td { border: 1px solid #ddd; padding: 10px; text-align: left; }
th { background-color: #4CAF50; color: white; }
tr:nth-child(even) { background-color: #f9f9f9; }
tr:hover { background-color: #f1f1f1; }
.state-finished { color: #4CAF50; font-weight: bold; }
.state-error { color: #f44336; font-weight: bold; }
.state-started { color: #2196F3; font-weight: bold; }
.state-queued { color: #ff9800; font-weight: bold; }
button {
margin: 5px;
padding: 10px 15px;
border: none;
border-radius: 4px;
cursor: pointer;
font-size: 14px;
transition: background 0.3s;
}
.btn-primary { background: #4CAF50; color: white; }
.btn-primary:hover { background: #45a049; }
.btn-secondary { background: #2196F3; color: white; }
.btn-secondary:hover { background: #0b7dda; }
.btn-danger { background: #f44336; color: white; }
.btn-danger:hover { background: #da190b; }
.status { margin-bottom: 20px; padding: 10px; background: #e7f3fe; border-left: 4px solid #2196F3; }
.status span { margin-left: 20px; color: #666; }
.filter-bar { margin: 10px 0; }
.filter-bar input { padding: 8px; width: 300px; border: 1px solid #ddd; border-radius: 4px; }
.stats { margin: 10px 0; font-size: 14px; color: #666; }
</style>
</head>
<body>
<div class="container">
<h1>🎬 AE Anons - Nexrender Job Manager</h1>
<div class="status">
<button class="btn-primary" onclick="generateJobs()">🔄 Generate New Jobs</button>
<button class="btn-danger" onclick="cleanupJobs()">🧹 Cleanup Finished Jobs</button>
<button class="btn-secondary" onclick="refreshJobs()">↻ Refresh</button>
<span id="statusMessage"></span>
</div>
<div class="filter-bar">
<input type="text" id="filterInput" placeholder="🔍 Filter by output file name..." onkeyup="filterTable()">
</div>
<div class="stats" id="stats"></div>
<table id="jobsTable">
<thead>
<tr>
<th>Output File</th>
<th>State</th>
<th>Created</th>
<th>Updated</th>
<th>UID</th>
</tr>
</thead>
<tbody id="jobsTableBody">
<tr><td colspan="5" style="text-align: center;">Loading...</td></tr>
</tbody>
</table>
</div>
<script>
let allJobs = [];
async function refreshJobs() {
try {
setStatus('Loading jobs...');
const response = await fetch('/api/jobs');
allJobs = await response.json();
renderJobs(allJobs);
updateStats();
setStatus(`Loaded ${allJobs.length} jobs`);
} catch (err) {
console.error(err);
setStatus('Error loading jobs: ' + err);
}
}
function renderJobs(jobs) {
const tbody = document.getElementById('jobsTableBody');
tbody.innerHTML = '';
if (jobs.length === 0) {
tbody.innerHTML = '<tr><td colspan="5" style="text-align: center;">No jobs found</td></tr>';
return;
}
jobs.sort((a, b) => {
// Сортируем по дате создания (новые сверху)
const dateA = a.created_at ? new Date(a.created_at) : new Date(0);
const dateB = b.created_at ? new Date(b.created_at) : new Date(0);
return dateB - dateA;
});
jobs.forEach(job => {
const row = tbody.insertRow();
// Output file (выделяем жирным)
const fileCell = row.insertCell();
fileCell.textContent = job.outfile_name;
fileCell.style.fontWeight = 'bold';
// State с цветовой кодировкой
const stateCell = row.insertCell();
stateCell.textContent = job.state;
stateCell.className = `state-${job.state}`;
// Форматируем даты
row.insertCell().textContent = formatDate(job.created_at);
row.insertCell().textContent = formatDate(job.updated_at);
// UID (обрезаем для компактности)
const uidCell = row.insertCell();
uidCell.textContent = job.uid.substring(0, 8) + '...';
uidCell.title = job.uid; // полный UID при наведении
});
}
function formatDate(dateStr) {
if (!dateStr) return '-';
try {
const date = new Date(dateStr);
return date.toLocaleString('ru-RU', {
day: '2-digit',
month: '2-digit',
hour: '2-digit',
minute: '2-digit'
});
} catch {
return dateStr;
}
}
function filterTable() {
const filter = document.getElementById('filterInput').value.toLowerCase();
const filtered = allJobs.filter(job =>
job.outfile_name.toLowerCase().includes(filter) ||
job.uid.toLowerCase().includes(filter)
);
renderJobs(filtered);
updateStats(filtered.length);
}
function updateStats(filteredCount = null) {
const total = allJobs.length;
const shown = filteredCount !== null ? filteredCount : total;
const states = allJobs.reduce((acc, job) => {
acc[job.state] = (acc[job.state] || 0) + 1;
return acc;
}, {});
const statsDiv = document.getElementById('stats');
const stateText = Object.entries(states)
.map(([state, count]) => `${state}: ${count}`)
.join(' | ');
statsDiv.innerHTML = `Total: ${total} jobs${filteredCount !== null ? ` | Showing: ${shown}` : ''} | ${stateText}`;
}
async function generateJobs() {
setStatus('⏳ Generating jobs...');
try {
const response = await fetch('/api/generate', { method: 'POST' });
if (response.ok) {
setStatus('✅ Job generation started. Will refresh in 10 seconds...');
setTimeout(() => refreshJobs(), 10000);
} else {
const text = await response.text();
setStatus('❌ Error: ' + text);
}
} catch (err) {
setStatus('❌ Error: ' + err);
}
}
async function cleanupJobs() {
setStatus('🧹 Cleaning up finished jobs...');
try {
const response = await fetch('/api/cleanup', { method: 'POST' });
if (response.ok) {
setStatus('✅ Cleanup completed. Refreshing...');
await refreshJobs();
} else {
const text = await response.text();
setStatus('❌ Error: ' + text);
}
} catch (err) {
setStatus('❌ Error: ' + err);
}
}
function setStatus(msg) {
document.getElementById('statusMessage').textContent = msg;
}
// Initial load
refreshJobs();
// Auto-refresh every 10 seconds
setInterval(refreshJobs, 10000);
</script>
</body>
</html>

186
src/web.rs Normal file
View File

@@ -0,0 +1,186 @@
use crate::config::Config;
use crate::processor::{cleanup_finished_jobs, fetch_all_jobs, process_spreadsheet};
use axum::{
extract::State,
http::StatusCode,
response::{Html, IntoResponse, Json},
routing::{get, post},
Router,
};
use serde::Serialize;
use std::net::SocketAddr;
use std::sync::Arc;
use tokio::sync::Mutex;
use tower_http::trace::TraceLayer;
#[derive(Clone)]
pub struct AppState {
pub config: Config,
pub last_generation: Arc<Mutex<Option<chrono::DateTime<chrono::Local>>>>,
}
#[derive(Serialize)]
pub struct JobInfo {
pub uid: String,
pub outfile_name: String,
pub state: String,
pub created_at: Option<String>,
pub updated_at: Option<String>,
}
impl JobInfo {
fn from_nexrender_json(job: &serde_json::Value) -> Option<Self> {
let uid = job.get("uid")?.as_str()?.to_string();
let state = job.get("state")?.as_str()?.to_string();
// Извлекаем имя выходного файла из postrender actions
let outfile_name = job
.get("actions")
.and_then(|a| a.get("postrender"))
.and_then(|p| p.as_array())
.and_then(|arr| {
// Ищем действие copy (в нём финальный путь)
arr.iter()
.find_map(|action| {
// Проверяем, что это действие copy
action
.get("module")
.and_then(|m| m.as_str())
.filter(|&m| m == "@nexrender/action-copy")
.and_then(|_| {
// Извлекаем output из copy
action.get("output").and_then(|o| o.as_str())
})
})
// Если copy не найдено, пробуем encode
.or_else(|| {
arr.iter()
.find_map(|action| action.get("output").and_then(|o| o.as_str()))
})
})
.map(|path| {
// Извлекаем только имя файла из полного пути
std::path::Path::new(path)
.file_name()
.and_then(|n| n.to_str())
.unwrap_or(path)
.to_string()
})
.unwrap_or_else(|| {
// Если не удалось извлечь, используем UID
format!("job_{}", uid)
});
Some(JobInfo {
uid,
outfile_name,
state,
created_at: job
.get("createdAt")
.and_then(|v| v.as_str())
.map(|s| s.to_string()),
updated_at: job
.get("updatedAt")
.and_then(|v| v.as_str())
.map(|s| s.to_string()),
})
}
}
pub async fn run_web_server(config: Config) -> anyhow::Result<()> {
let state = AppState {
config,
last_generation: Arc::new(Mutex::new(None)),
};
let app = Router::new()
.route("/", get(index_page))
.route("/api/jobs", get(list_jobs))
.route("/api/generate", post(generate_jobs))
.route("/api/cleanup", post(cleanup_jobs))
.route("/api/status", get(get_status))
.layer(TraceLayer::new_for_http())
.with_state(state);
let addr: SocketAddr = "0.0.0.0:3000".parse()?;
log::info!("Web server listening on http://{}", addr);
let listener = tokio::net::TcpListener::bind(addr).await?;
axum::serve(listener, app).await?;
Ok(())
}
// Остальные обработчики остаются без изменений...
async fn index_page() -> Html<&'static str> {
Html(include_str!("static/index.html"))
}
async fn list_jobs(State(state): State<AppState>) -> Result<Json<Vec<JobInfo>>, AppError> {
let jobs_json = fetch_all_jobs(&state.config.nexrender_api_url)
.await
.map_err(|e| AppError(StatusCode::INTERNAL_SERVER_ERROR, e.to_string()))?;
let jobs: Vec<JobInfo> = jobs_json
.iter()
.filter_map(JobInfo::from_nexrender_json)
.collect();
Ok(Json(jobs))
}
async fn generate_jobs(State(state): State<AppState>) -> Result<impl IntoResponse, AppError> {
let mut last_gen = state.last_generation.lock().await;
if let Some(last) = *last_gen {
let elapsed = chrono::Local::now().signed_duration_since(last);
if elapsed.num_seconds() < 5 {
return Err(AppError(
StatusCode::TOO_MANY_REQUESTS,
"Generation already in progress or too recent".to_string(),
));
}
}
*last_gen = Some(chrono::Local::now());
drop(last_gen);
let config = state.config.clone();
let last_gen_clone = state.last_generation.clone();
tokio::spawn(async move {
match process_spreadsheet(&config).await {
Ok(submitted) => {
log::info!("Generation completed, {} jobs submitted", submitted.len());
}
Err(e) => {
log::error!("Generation failed: {}", e);
}
}
*last_gen_clone.lock().await = None;
});
Ok((StatusCode::ACCEPTED, "Job generation started"))
}
async fn cleanup_jobs(State(state): State<AppState>) -> Result<impl IntoResponse, AppError> {
cleanup_finished_jobs(&state.config.nexrender_api_url)
.await
.map_err(|e| AppError(StatusCode::INTERNAL_SERVER_ERROR, e.to_string()))?;
Ok((StatusCode::OK, "Cleanup completed"))
}
async fn get_status(State(state): State<AppState>) -> Result<Json<serde_json::Value>, AppError> {
let last_gen = *state.last_generation.lock().await;
let status = serde_json::json!({
"last_generation": last_gen.map(|dt| dt.to_rfc3339()),
"nexrender_api": state.config.nexrender_api_url,
});
Ok(Json(status))
}
// Error handling
struct AppError(StatusCode, String);
impl IntoResponse for AppError {
fn into_response(self) -> axum::response::Response {
(self.0, self.1).into_response()
}
}