diff --git a/src/config.rs b/src/config.rs index d3cf395..e86d740 100644 --- a/src/config.rs +++ b/src/config.rs @@ -23,7 +23,6 @@ pub struct Config { pub smb_upload_url: Option, pub smb_username: Option, pub smb_password: Option, - pub smb_workgroup: Option, // SMB source folder (where videos are stored locally before upload) pub smb_source_folder: String, // Database @@ -61,10 +60,9 @@ impl Config { smb_upload_url: env::var("SMB_UPLOAD_URL").ok(), smb_username: env::var("SMB_USERNAME").ok(), smb_password: env::var("SMB_PASSWORD").ok(), - smb_workgroup: env::var("SMB_WORKGROUP").ok(), smb_source_folder, database_url: env::var("DATABASE_URL") .unwrap_or_else(|_| "sqlite:ae_anons.db".to_string()), }) } -} +} \ No newline at end of file diff --git a/src/storage.rs b/src/storage.rs index 1fe6ebd..3afb0e7 100644 --- a/src/storage.rs +++ b/src/storage.rs @@ -15,50 +15,59 @@ pub struct JobRecord { } pub async fn init_db(database_url: &str) -> Result { - match SqlitePool::connect(database_url).await { - Ok(pool) => { - // Создаём таблицы - sqlx::query( - r#" - CREATE TABLE IF NOT EXISTS jobs ( - uid TEXT PRIMARY KEY, - filename TEXT NOT NULL, - status TEXT NOT NULL, - output_path TEXT NOT NULL, - approved BOOLEAN NOT NULL DEFAULT 0, - created_at TIMESTAMP NOT NULL, - updated_at TIMESTAMP NOT NULL - ) - "#, + // Для in-memory БД (тестирование) + if database_url == "sqlite::memory:" { + let pool = SqlitePool::connect("sqlite::memory:").await?; + sqlx::query( + r#" + CREATE TABLE IF NOT EXISTS jobs ( + uid TEXT PRIMARY KEY, + filename TEXT NOT NULL, + status TEXT NOT NULL, + output_path TEXT NOT NULL, + approved BOOLEAN NOT NULL DEFAULT 0, + created_at TIMESTAMP NOT NULL, + updated_at TIMESTAMP NOT NULL ) - .execute(&pool) - .await?; - Ok(pool) - } - Err(e) => { - log::warn!("Failed to open database file '{}': {}", database_url, e); - log::warn!("Falling back to in-memory database. Data will not be persisted!"); - let pool = SqlitePool::connect("sqlite::memory:").await?; - sqlx::query( - r#" - CREATE TABLE IF NOT EXISTS jobs ( - uid TEXT PRIMARY KEY, - filename TEXT NOT NULL, - status TEXT NOT NULL, - output_path TEXT NOT NULL, - approved BOOLEAN NOT NULL DEFAULT 0, - created_at TIMESTAMP NOT NULL, - updated_at TIMESTAMP NOT NULL - ) - "#, - ) - .execute(&pool) - .await?; - Ok(pool) + "#, + ) + .execute(&pool) + .await?; + return Ok(pool); + } + + // Для файловой БД + let path = database_url.strip_prefix("sqlite:").unwrap_or(database_url); + // Создаём родительскую директорию, если её нет + if let Some(parent) = std::path::Path::new(path).parent() { + if !parent.exists() { + tokio::fs::create_dir_all(parent).await + .map_err(|e| anyhow::anyhow!("Failed to create database directory '{}': {}", parent.display(), e))?; } } + // Подключаемся (если файла нет, SQLite создаст его автоматически) + let pool = SqlitePool::connect(database_url).await + .map_err(|e| anyhow::anyhow!("Failed to open database '{}': {}", database_url, e))?; + // Создаём таблицу, если её нет + sqlx::query( + r#" + CREATE TABLE IF NOT EXISTS jobs ( + uid TEXT PRIMARY KEY, + filename TEXT NOT NULL, + status TEXT NOT NULL, + output_path TEXT NOT NULL, + approved BOOLEAN NOT NULL DEFAULT 0, + created_at TIMESTAMP NOT NULL, + updated_at TIMESTAMP NOT NULL + ) + "#, + ) + .execute(&pool) + .await?; + Ok(pool) } + pub async fn upsert_job(pool: &SqlitePool, record: &JobRecord) -> Result<()> { sqlx::query( r#" diff --git a/src/web.rs b/src/web.rs index a522226..a69b0bd 100644 --- a/src/web.rs +++ b/src/web.rs @@ -19,6 +19,7 @@ use tokio::sync::broadcast; use tokio_util::io::ReaderStream; use tower_http::trace::TraceLayer; use url::Url; +use chrono::Utc; #[derive(Clone)] pub struct AppState { @@ -83,24 +84,45 @@ async fn sync_jobs_status(state: AppState) { let mut interval = tokio::time::interval(std::time::Duration::from_secs(5)); loop { interval.tick().await; - if let Ok(jobs_json) = fetch_all_jobs(&state.config.nexrender_api_url).await { - for job_json in jobs_json { - if let Some(job) = JobInfo::from_nexrender_json(&job_json) { - // Получаем текущую запись из БД - if let Ok(Some(mut record)) = storage::get_job(&state.db, &job.uid).await { - if record.status != job.state { - // Статус изменился – обновляем БД и шлём событие - record.status = job.state.clone(); - record.updated_at = chrono::Utc::now(); - if let Err(e) = storage::upsert_job(&state.db, &record).await { - log::error!("Failed to update job status in DB: {}", e); - } else { - let _ = state.ws_tx.send(WsEvent::JobUpdated(record)); + log::debug!("🔄 Syncing jobs from Nexrender..."); + match fetch_all_jobs(&state.config.nexrender_api_url).await { + Ok(jobs_json) => { + log::info!("📊 Fetched {} jobs from Nexrender", jobs_json.len()); + for job_json in jobs_json { + if let Some(job) = JobInfo::from_nexrender_json(&job_json) { + log::debug!("🔍 Job from Nexrender: uid={}, state={}", job.uid, job.state); + match storage::get_job(&state.db, &job.uid).await { + Ok(Some(record)) => { + if record.status != job.state { + log::info!("🔄 Status changed for job {}: {} -> {}", job.uid, record.status, job.state); + let mut updated_record = record; + updated_record.status = job.state.clone(); + updated_record.updated_at = Utc::now(); + if let Err(e) = storage::upsert_job(&state.db, &updated_record).await { + log::error!("❌ Failed to update job status in DB: {}", e); + } else { + let _ = state.ws_tx.send(WsEvent::JobUpdated(updated_record)); + log::debug!("📨 Sent WebSocket event for job {}", job.uid); + } + } else { + log::debug!("✅ No change for job {}", job.uid); + } + } + Ok(None) => { + log::warn!("⚠️ Job {} not found in DB, skipping", job.uid); + } + Err(e) => { + log::error!("❌ DB error for job {}: {}", job.uid, e); } } + } else { + log::warn!("⚠️ Failed to parse job from JSON: {:?}", job_json); } } } + Err(e) => { + log::error!("❌ Failed to fetch jobs from Nexrender: {}", e); + } } } } @@ -308,7 +330,6 @@ async fn approve_job( "SMB_PASSWORD not configured".to_string(), ) })?; - let _smb_workgroup = state.config.smb_workgroup.clone().unwrap_or_default(); // *** ИСПРАВЛЕНИЕ ЗДЕСЬ *** let src_path = std::path::Path::new(&state.config.smb_source_folder).join(&record.filename); @@ -319,27 +340,40 @@ async fn approve_job( )); } + // Исправление для SMB URL обработки let url = Url::parse(smb_url) .map_err(|e| AppError(StatusCode::BAD_REQUEST, format!("Invalid SMB URL: {}", e)))?; + let server = url .host_str() .ok_or_else(|| AppError(StatusCode::BAD_REQUEST, "No host in SMB URL".to_string()))?; - let share = url - .path() - .trim_start_matches('/') - .split('/') - .next() - .unwrap_or(""); - if share.is_empty() { + + // Исправленный способ извлечения share и path + let path_parts: Vec<&str> = url.path().trim_matches('/').split('/').collect(); + if path_parts.is_empty() { return Err(AppError( StatusCode::BAD_REQUEST, "No share in SMB URL".to_string(), )); } + + // Первый элемент пути - это имя шары + let share = path_parts[0]; + + // Остальная часть пути будет использоваться как путь к файлу (если есть) + let remote_path = if path_parts.len() > 1 { + path_parts.iter().skip(1).map(|&s| s.to_string()).collect::>().join("/") + } else { + String::new() + }; + // Исправленное подключение к SMB let client = Client::new(ClientConfig::default()); + + // Подключаемся к шаре let target_path = UncPath::from_str(&format!("\\\\{}\\{}", server, share)) .map_err(|e| AppError(StatusCode::BAD_REQUEST, format!("Invalid UNC path: {}", e)))?; + client .share_connect(&target_path, smb_user, smb_pass.clone()) .await @@ -350,6 +384,7 @@ async fn approve_job( ) })?; + // Читаем содержимое файла let data = tokio::fs::read(&src_path).await.map_err(|e| { AppError( StatusCode::INTERNAL_SERVER_ERROR, @@ -357,13 +392,14 @@ async fn approve_job( ) })?; - let remote_path = url.path().trim_start_matches('/'); - let remote_filename = if remote_path.is_empty() { - record.filename.clone() - } else { + // Определяем путь для сохранения на SMB + let remote_filename = if !remote_path.is_empty() { format!("{}/{}", remote_path, record.filename) + } else { + record.filename.clone() }; + // Создаем полный UNC путь к файлу в шаре let file_to_open = target_path.with_path(&remote_filename); let file_open_args = FileCreateArgs::make_overwrite(FileAttributes::default(), CreateOptions::default()); @@ -376,14 +412,17 @@ async fn approve_job( format!("Create file failed: {}", e), ) })?; + let remote_file = resource.unwrap_file(); + // Записываем данные на SMB remote_file.write_at(&data, 0).await.map_err(|e| { AppError( StatusCode::INTERNAL_SERVER_ERROR, format!("Write to SMB failed: {}", e), ) })?; + remote_file.close().await.map_err(|e| { AppError( StatusCode::INTERNAL_SERVER_ERROR, @@ -391,12 +430,15 @@ async fn approve_job( ) })?; + // Обновляем статус в БД storage::approve_job(&state.db, &uid) .await .map_err(|e| AppError(StatusCode::INTERNAL_SERVER_ERROR, e.to_string()))?; record.approved = true; - let _ = state.ws_tx.send(WsEvent::JobUpdated(record)); + + // Отправляем обновление через WebSocket + let _ = state.ws_tx.send(WsEvent::JobUpdated(record.clone())); Ok(Json(json!({"status": "approved"}))) } @@ -423,7 +465,6 @@ async fn handle_socket(mut socket: axum::extract::ws::WebSocket, state: AppState #[derive(Debug, Clone)] struct JobInfo { uid: String, - outfile_name: String, state: String, } @@ -432,36 +473,8 @@ impl JobInfo { let uid = job.get("uid")?.as_str()?.to_string(); let state = job.get("state")?.as_str()?.to_string(); - let outfile_name = job - .get("actions") - .and_then(|a| a.get("postrender")) - .and_then(|p| p.as_array()) - .and_then(|arr| { - arr.iter() - .find_map(|action| { - action - .get("module") - .and_then(|m| m.as_str()) - .filter(|&m| m == "@nexrender/action-copy") - .and_then(|_| action.get("output").and_then(|o| o.as_str())) - }) - .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(|| format!("job_{}", uid)); - Some(JobInfo { uid, - outfile_name, state, }) } @@ -473,4 +486,4 @@ impl IntoResponse for AppError { fn into_response(self) -> Response { (self.0, self.1).into_response() } -} +} \ No newline at end of file