Всё ещё не работает

This commit is contained in:
2026-05-06 12:03:25 +03:00
parent 02fd7a65a7
commit 6fcc518438
3 changed files with 118 additions and 98 deletions

View File

@@ -23,7 +23,6 @@ pub struct Config {
pub smb_upload_url: Option<String>, pub smb_upload_url: Option<String>,
pub smb_username: Option<String>, pub smb_username: Option<String>,
pub smb_password: Option<String>, pub smb_password: Option<String>,
pub smb_workgroup: Option<String>,
// SMB source folder (where videos are stored locally before upload) // SMB source folder (where videos are stored locally before upload)
pub smb_source_folder: String, pub smb_source_folder: String,
// Database // Database
@@ -61,7 +60,6 @@ impl Config {
smb_upload_url: env::var("SMB_UPLOAD_URL").ok(), smb_upload_url: env::var("SMB_UPLOAD_URL").ok(),
smb_username: env::var("SMB_USERNAME").ok(), smb_username: env::var("SMB_USERNAME").ok(),
smb_password: env::var("SMB_PASSWORD").ok(), smb_password: env::var("SMB_PASSWORD").ok(),
smb_workgroup: env::var("SMB_WORKGROUP").ok(),
smb_source_folder, smb_source_folder,
database_url: env::var("DATABASE_URL") database_url: env::var("DATABASE_URL")
.unwrap_or_else(|_| "sqlite:ae_anons.db".to_string()), .unwrap_or_else(|_| "sqlite:ae_anons.db".to_string()),

View File

@@ -15,29 +15,8 @@ pub struct JobRecord {
} }
pub async fn init_db(database_url: &str) -> Result<SqlitePool> { pub async fn init_db(database_url: &str) -> Result<SqlitePool> {
match SqlitePool::connect(database_url).await { // Для in-memory БД (тестирование)
Ok(pool) => { if database_url == "sqlite::memory:" {
// Создаём таблицы
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?; let pool = SqlitePool::connect("sqlite::memory:").await?;
sqlx::query( sqlx::query(
r#" r#"
@@ -54,11 +33,41 @@ pub async fn init_db(database_url: &str) -> Result<SqlitePool> {
) )
.execute(&pool) .execute(&pool)
.await?; .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) Ok(pool)
}
}
} }
pub async fn upsert_job(pool: &SqlitePool, record: &JobRecord) -> Result<()> { pub async fn upsert_job(pool: &SqlitePool, record: &JobRecord) -> Result<()> {
sqlx::query( sqlx::query(
r#" r#"

View File

@@ -19,6 +19,7 @@ use tokio::sync::broadcast;
use tokio_util::io::ReaderStream; use tokio_util::io::ReaderStream;
use tower_http::trace::TraceLayer; use tower_http::trace::TraceLayer;
use url::Url; use url::Url;
use chrono::Utc;
#[derive(Clone)] #[derive(Clone)]
pub struct AppState { pub struct AppState {
@@ -83,23 +84,44 @@ async fn sync_jobs_status(state: AppState) {
let mut interval = tokio::time::interval(std::time::Duration::from_secs(5)); let mut interval = tokio::time::interval(std::time::Duration::from_secs(5));
loop { loop {
interval.tick().await; interval.tick().await;
if let Ok(jobs_json) = fetch_all_jobs(&state.config.nexrender_api_url).await { 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 { for job_json in jobs_json {
if let Some(job) = JobInfo::from_nexrender_json(&job_json) { if let Some(job) = JobInfo::from_nexrender_json(&job_json) {
// Получаем текущую запись из БД log::debug!("🔍 Job from Nexrender: uid={}, state={}", job.uid, job.state);
if let Ok(Some(mut record)) = storage::get_job(&state.db, &job.uid).await { match storage::get_job(&state.db, &job.uid).await {
Ok(Some(record)) => {
if record.status != job.state { if record.status != job.state {
// Статус изменился обновляем БД и шлём событие log::info!("🔄 Status changed for job {}: {} -> {}", job.uid, record.status, job.state);
record.status = job.state.clone(); let mut updated_record = record;
record.updated_at = chrono::Utc::now(); updated_record.status = job.state.clone();
if let Err(e) = storage::upsert_job(&state.db, &record).await { updated_record.updated_at = Utc::now();
log::error!("Failed to update job status in DB: {}", e); if let Err(e) = storage::upsert_job(&state.db, &updated_record).await {
log::error!("❌ Failed to update job status in DB: {}", e);
} else { } else {
let _ = state.ws_tx.send(WsEvent::JobUpdated(record)); 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(), "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); 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) let url = Url::parse(smb_url)
.map_err(|e| AppError(StatusCode::BAD_REQUEST, format!("Invalid SMB URL: {}", e)))?; .map_err(|e| AppError(StatusCode::BAD_REQUEST, format!("Invalid SMB URL: {}", e)))?;
let server = url let server = url
.host_str() .host_str()
.ok_or_else(|| AppError(StatusCode::BAD_REQUEST, "No host in SMB URL".to_string()))?; .ok_or_else(|| AppError(StatusCode::BAD_REQUEST, "No host in SMB URL".to_string()))?;
let share = url
.path() // Исправленный способ извлечения share и path
.trim_start_matches('/') let path_parts: Vec<&str> = url.path().trim_matches('/').split('/').collect();
.split('/') if path_parts.is_empty() {
.next()
.unwrap_or("");
if share.is_empty() {
return Err(AppError( return Err(AppError(
StatusCode::BAD_REQUEST, StatusCode::BAD_REQUEST,
"No share in SMB URL".to_string(), "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::<Vec<_>>().join("/")
} else {
String::new()
};
// Исправленное подключение к SMB
let client = Client::new(ClientConfig::default()); let client = Client::new(ClientConfig::default());
// Подключаемся к шаре
let target_path = UncPath::from_str(&format!("\\\\{}\\{}", server, share)) let target_path = UncPath::from_str(&format!("\\\\{}\\{}", server, share))
.map_err(|e| AppError(StatusCode::BAD_REQUEST, format!("Invalid UNC path: {}", e)))?; .map_err(|e| AppError(StatusCode::BAD_REQUEST, format!("Invalid UNC path: {}", e)))?;
client client
.share_connect(&target_path, smb_user, smb_pass.clone()) .share_connect(&target_path, smb_user, smb_pass.clone())
.await .await
@@ -350,6 +384,7 @@ async fn approve_job(
) )
})?; })?;
// Читаем содержимое файла
let data = tokio::fs::read(&src_path).await.map_err(|e| { let data = tokio::fs::read(&src_path).await.map_err(|e| {
AppError( AppError(
StatusCode::INTERNAL_SERVER_ERROR, StatusCode::INTERNAL_SERVER_ERROR,
@@ -357,13 +392,14 @@ async fn approve_job(
) )
})?; })?;
let remote_path = url.path().trim_start_matches('/'); // Определяем путь для сохранения на SMB
let remote_filename = if remote_path.is_empty() { let remote_filename = if !remote_path.is_empty() {
record.filename.clone()
} else {
format!("{}/{}", remote_path, record.filename) format!("{}/{}", remote_path, record.filename)
} else {
record.filename.clone()
}; };
// Создаем полный UNC путь к файлу в шаре
let file_to_open = target_path.with_path(&remote_filename); let file_to_open = target_path.with_path(&remote_filename);
let file_open_args = let file_open_args =
FileCreateArgs::make_overwrite(FileAttributes::default(), CreateOptions::default()); FileCreateArgs::make_overwrite(FileAttributes::default(), CreateOptions::default());
@@ -376,14 +412,17 @@ async fn approve_job(
format!("Create file failed: {}", e), format!("Create file failed: {}", e),
) )
})?; })?;
let remote_file = resource.unwrap_file(); let remote_file = resource.unwrap_file();
// Записываем данные на SMB
remote_file.write_at(&data, 0).await.map_err(|e| { remote_file.write_at(&data, 0).await.map_err(|e| {
AppError( AppError(
StatusCode::INTERNAL_SERVER_ERROR, StatusCode::INTERNAL_SERVER_ERROR,
format!("Write to SMB failed: {}", e), format!("Write to SMB failed: {}", e),
) )
})?; })?;
remote_file.close().await.map_err(|e| { remote_file.close().await.map_err(|e| {
AppError( AppError(
StatusCode::INTERNAL_SERVER_ERROR, StatusCode::INTERNAL_SERVER_ERROR,
@@ -391,12 +430,15 @@ async fn approve_job(
) )
})?; })?;
// Обновляем статус в БД
storage::approve_job(&state.db, &uid) storage::approve_job(&state.db, &uid)
.await .await
.map_err(|e| AppError(StatusCode::INTERNAL_SERVER_ERROR, e.to_string()))?; .map_err(|e| AppError(StatusCode::INTERNAL_SERVER_ERROR, e.to_string()))?;
record.approved = true; 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"}))) Ok(Json(json!({"status": "approved"})))
} }
@@ -423,7 +465,6 @@ async fn handle_socket(mut socket: axum::extract::ws::WebSocket, state: AppState
#[derive(Debug, Clone)] #[derive(Debug, Clone)]
struct JobInfo { struct JobInfo {
uid: String, uid: String,
outfile_name: String,
state: String, state: String,
} }
@@ -432,36 +473,8 @@ impl JobInfo {
let uid = job.get("uid")?.as_str()?.to_string(); let uid = job.get("uid")?.as_str()?.to_string();
let state = job.get("state")?.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 { Some(JobInfo {
uid, uid,
outfile_name,
state, state,
}) })
} }