55bca92691
Core features: - Class 1 T.30 protocol: full send/receive implementation - HDLC: DLE-stuffing, FCS strip, USR5637 bit-reversal handling - T.4 MH encoder/decoder (1728px A4 standard) - Document pipeline: PDF (Ghostscript), PNG, TIFF input - Width clamping: US Letter 1734px → 1728px fax standard - Cover page: CJK rasterization (TW/CN/JP/EN), TIFF + HTML output - OCR verification: Tesseract 5 with eng+chi_tra, CJK space-tolerant - API server (axum): health, send, jobs, cover, retry, cancel - Background worker: auto-poll queue, speed fallback, retry policy - Modem detection, pool management Real-world test results (2026-07-23): - V90 → 25153038: 4 pages, V.17 12000 bps, 2:33 ✅ - USR5637 → 25153038: 4 pages, V.17 12000 bps, 2:26 ✅ - Both faxes confirmed received on remote machine Tested: loopback (100% pixel match), multi-page, all input formats, cover pages, OCR verify, API endpoints, worker processing. 13 unit tests pass, 0 new clippy warnings.
375 lines
10 KiB
Rust
375 lines
10 KiB
Rust
use axum::{
|
|
extract::{Path, State},
|
|
http::StatusCode,
|
|
response::Json,
|
|
routing::{delete, get, post, put},
|
|
Router,
|
|
};
|
|
use serde::{Deserialize, Serialize};
|
|
use std::sync::Arc;
|
|
use tokio::sync::Mutex;
|
|
use tracing::info;
|
|
|
|
use crate::config_new::AppConfig;
|
|
use crate::error::Result as FaxResult;
|
|
use crate::queue::{FaxJob, FaxQueue, JobId, JobStatus};
|
|
use crate::worker::FaxWorker;
|
|
|
|
#[derive(Clone)]
|
|
pub struct AppState {
|
|
pub config: AppConfig,
|
|
pub queue: Arc<Mutex<FaxQueue>>,
|
|
pub worker: Arc<Mutex<Option<FaxWorker>>>,
|
|
}
|
|
|
|
#[derive(Serialize)]
|
|
pub struct HealthResponse {
|
|
pub status: String,
|
|
pub uptime_seconds: u64,
|
|
pub modems: Vec<ModemHealthInfo>,
|
|
pub queue: QueueHealthInfo,
|
|
}
|
|
|
|
#[derive(Serialize)]
|
|
pub struct ModemHealthInfo {
|
|
pub name: String,
|
|
pub device: String,
|
|
pub status: String,
|
|
}
|
|
|
|
#[derive(Serialize)]
|
|
pub struct QueueHealthInfo {
|
|
pub pending: usize,
|
|
pub active: usize,
|
|
pub failed: usize,
|
|
}
|
|
|
|
#[derive(Deserialize)]
|
|
pub struct SendFaxRequest {
|
|
pub recipient: String,
|
|
pub document_path: String,
|
|
#[serde(default)]
|
|
pub cover_to: Option<String>,
|
|
#[serde(default)]
|
|
pub cover_from: Option<String>,
|
|
#[serde(default)]
|
|
pub cover_subject: Option<String>,
|
|
#[serde(default)]
|
|
pub cover_notes: Option<String>,
|
|
}
|
|
|
|
#[derive(Serialize)]
|
|
pub struct SendFaxResponse {
|
|
pub job_id: JobId,
|
|
pub message: String,
|
|
}
|
|
|
|
#[derive(Deserialize)]
|
|
pub struct UpdateCoverRequest {
|
|
pub to: Option<String>,
|
|
pub from: Option<String>,
|
|
pub subject: Option<String>,
|
|
pub notes: Option<String>,
|
|
}
|
|
|
|
#[derive(Serialize)]
|
|
pub struct UpdateCoverResponse {
|
|
pub job_id: JobId,
|
|
pub cover_to: Option<String>,
|
|
pub cover_from: Option<String>,
|
|
pub cover_subject: Option<String>,
|
|
pub cover_notes: Option<String>,
|
|
}
|
|
|
|
#[derive(Serialize)]
|
|
pub struct JobResponse {
|
|
pub id: JobId,
|
|
pub recipient: String,
|
|
pub status: String,
|
|
pub pages: u32,
|
|
pub retries: u8,
|
|
pub created_at: String,
|
|
pub updated_at: String,
|
|
pub cover_to: Option<String>,
|
|
pub cover_from: Option<String>,
|
|
pub cover_subject: Option<String>,
|
|
pub cover_notes: Option<String>,
|
|
}
|
|
|
|
#[derive(Serialize)]
|
|
pub struct JobListResponse {
|
|
pub jobs: Vec<JobResponse>,
|
|
pub total: usize,
|
|
}
|
|
|
|
#[derive(Serialize)]
|
|
pub struct ErrorResponse {
|
|
pub error: String,
|
|
pub message: String,
|
|
}
|
|
|
|
async fn health(State(state): State<AppState>) -> Json<HealthResponse> {
|
|
let modems: Vec<ModemHealthInfo> = state.config.modems.iter().map(|m| {
|
|
ModemHealthInfo {
|
|
name: m.name.clone(),
|
|
device: m.device.clone(),
|
|
status: "idle".to_string(),
|
|
}
|
|
}).collect();
|
|
|
|
let queue = state.queue.lock().await;
|
|
let jobs = queue.list();
|
|
let pending = jobs.iter().filter(|j| matches!(j.status, JobStatus::Queued)).count();
|
|
let failed = jobs.iter().filter(|j| matches!(j.status, JobStatus::Failed(_))).count();
|
|
|
|
Json(HealthResponse {
|
|
status: "healthy".to_string(),
|
|
uptime_seconds: 0,
|
|
modems,
|
|
queue: QueueHealthInfo {
|
|
pending,
|
|
active: 0,
|
|
failed,
|
|
},
|
|
})
|
|
}
|
|
|
|
async fn send_fax(
|
|
State(state): State<AppState>,
|
|
Json(req): Json<SendFaxRequest>,
|
|
) -> Result<Json<SendFaxResponse>, (StatusCode, Json<ErrorResponse>)> {
|
|
let mut job = FaxJob::new(req.recipient, req.document_path);
|
|
job.cover_to = req.cover_to;
|
|
job.cover_from = req.cover_from;
|
|
job.cover_subject = req.cover_subject;
|
|
job.cover_notes = req.cover_notes;
|
|
let job_id = job.id;
|
|
|
|
let mut queue = state.queue.lock().await;
|
|
queue.enqueue(job).map_err(|e| {
|
|
(
|
|
StatusCode::INTERNAL_SERVER_ERROR,
|
|
Json(ErrorResponse {
|
|
error: "queue_error".to_string(),
|
|
message: e.to_string(),
|
|
}),
|
|
)
|
|
})?;
|
|
|
|
info!(job_id = %job_id, "Fax job queued via API");
|
|
|
|
Ok(Json(SendFaxResponse {
|
|
job_id,
|
|
message: "Job queued successfully".to_string(),
|
|
}))
|
|
}
|
|
|
|
async fn list_jobs(State(state): State<AppState>) -> Json<JobListResponse> {
|
|
let queue = state.queue.lock().await;
|
|
let jobs = queue.list();
|
|
let total = jobs.len();
|
|
|
|
let responses: Vec<JobResponse> = jobs
|
|
.into_iter()
|
|
.map(|j| JobResponse {
|
|
id: j.id,
|
|
recipient: j.recipient,
|
|
status: format!("{:?}", j.status),
|
|
pages: j.pages,
|
|
retries: j.retries,
|
|
created_at: j.created_at.to_rfc3339(),
|
|
updated_at: j.updated_at.to_rfc3339(),
|
|
cover_to: j.cover_to,
|
|
cover_from: j.cover_from,
|
|
cover_subject: j.cover_subject,
|
|
cover_notes: j.cover_notes,
|
|
})
|
|
.collect();
|
|
|
|
Json(JobListResponse { jobs: responses, total })
|
|
}
|
|
|
|
async fn get_job(
|
|
State(state): State<AppState>,
|
|
Path(id): Path<JobId>,
|
|
) -> Result<Json<JobResponse>, (StatusCode, Json<ErrorResponse>)> {
|
|
let queue = state.queue.lock().await;
|
|
queue.get(&id).map(|j| Json(JobResponse {
|
|
id: j.id,
|
|
recipient: j.recipient,
|
|
status: format!("{:?}", j.status),
|
|
pages: j.pages,
|
|
retries: j.retries,
|
|
created_at: j.created_at.to_rfc3339(),
|
|
updated_at: j.updated_at.to_rfc3339(),
|
|
cover_to: j.cover_to,
|
|
cover_from: j.cover_from,
|
|
cover_subject: j.cover_subject,
|
|
cover_notes: j.cover_notes,
|
|
})).ok_or((
|
|
StatusCode::NOT_FOUND,
|
|
Json(ErrorResponse {
|
|
error: "not_found".to_string(),
|
|
message: format!("Job {} not found", id),
|
|
}),
|
|
))
|
|
}
|
|
|
|
async fn cancel_job(
|
|
State(state): State<AppState>,
|
|
Path(id): Path<JobId>,
|
|
) -> Result<StatusCode, (StatusCode, Json<ErrorResponse>)> {
|
|
let mut queue = state.queue.lock().await;
|
|
|
|
if queue.get(&id).is_none() {
|
|
return Err((
|
|
StatusCode::NOT_FOUND,
|
|
Json(ErrorResponse {
|
|
error: "not_found".to_string(),
|
|
message: format!("Job {} not found", id),
|
|
}),
|
|
));
|
|
}
|
|
|
|
queue.update_status(&id, JobStatus::Cancelled).map_err(|e| {
|
|
(
|
|
StatusCode::INTERNAL_SERVER_ERROR,
|
|
Json(ErrorResponse {
|
|
error: "update_failed".to_string(),
|
|
message: e.to_string(),
|
|
}),
|
|
)
|
|
})?;
|
|
|
|
info!(job_id = %id, "Job cancelled via API");
|
|
Ok(StatusCode::NO_CONTENT)
|
|
}
|
|
|
|
async fn retry_job(
|
|
State(state): State<AppState>,
|
|
Path(id): Path<JobId>,
|
|
) -> Result<StatusCode, (StatusCode, Json<ErrorResponse>)> {
|
|
let mut queue = state.queue.lock().await;
|
|
|
|
if queue.get(&id).is_none() {
|
|
return Err((
|
|
StatusCode::NOT_FOUND,
|
|
Json(ErrorResponse {
|
|
error: "not_found".to_string(),
|
|
message: format!("Job {} not found", id),
|
|
}),
|
|
));
|
|
}
|
|
|
|
queue.update_status(&id, JobStatus::Queued).map_err(|e| {
|
|
(
|
|
StatusCode::INTERNAL_SERVER_ERROR,
|
|
Json(ErrorResponse {
|
|
error: "update_failed".to_string(),
|
|
message: e.to_string(),
|
|
}),
|
|
)
|
|
})?;
|
|
|
|
info!(job_id = %id, "Job retry requested via API");
|
|
Ok(StatusCode::OK)
|
|
}
|
|
|
|
async fn update_job_cover(
|
|
State(state): State<AppState>,
|
|
Path(id): Path<JobId>,
|
|
Json(req): Json<UpdateCoverRequest>,
|
|
) -> Result<Json<UpdateCoverResponse>, (StatusCode, Json<ErrorResponse>)> {
|
|
let mut queue = state.queue.lock().await;
|
|
|
|
queue
|
|
.update_cover(&id, req.to.clone(), req.from.clone(), req.subject.clone(), req.notes.clone())
|
|
.map_err(|e| {
|
|
(
|
|
StatusCode::INTERNAL_SERVER_ERROR,
|
|
Json(ErrorResponse {
|
|
error: "update_failed".to_string(),
|
|
message: e.to_string(),
|
|
}),
|
|
)
|
|
})?;
|
|
|
|
let job = queue.get(&id).ok_or((
|
|
StatusCode::NOT_FOUND,
|
|
Json(ErrorResponse {
|
|
error: "not_found".to_string(),
|
|
message: format!("Job {} not found", id),
|
|
}),
|
|
))?;
|
|
|
|
Ok(Json(UpdateCoverResponse {
|
|
job_id: job.id,
|
|
cover_to: job.cover_to,
|
|
cover_from: job.cover_from,
|
|
cover_subject: job.cover_subject,
|
|
cover_notes: job.cover_notes,
|
|
}))
|
|
}
|
|
|
|
pub fn create_router(state: AppState) -> Router {
|
|
Router::new()
|
|
.route("/api/health", get(health))
|
|
.route("/api/fax/send", post(send_fax))
|
|
.route("/api/fax/jobs", get(list_jobs))
|
|
.route("/api/fax/jobs/{id}", get(get_job))
|
|
.route("/api/fax/jobs/{id}", delete(cancel_job))
|
|
.route("/api/fax/jobs/{id}/retry", post(retry_job))
|
|
.route("/api/fax/jobs/{id}/cover", put(update_job_cover))
|
|
.with_state(state)
|
|
}
|
|
|
|
pub async fn start_server(config: AppConfig, queue: Arc<Mutex<FaxQueue>>) -> FaxResult<()> {
|
|
let fax_config = if let Some(modem) = config.primary_modem() {
|
|
crate::config::FaxConfig {
|
|
device: modem.device.clone(),
|
|
baud_rate: 115200,
|
|
station_id: config.fax.station_id.clone(),
|
|
header: config.fax.header.clone(),
|
|
..Default::default()
|
|
}
|
|
} else {
|
|
crate::config::FaxConfig::default()
|
|
};
|
|
|
|
let worker = crate::worker::FaxWorker::new(fax_config, queue.clone());
|
|
let worker_handle = Arc::new(Mutex::new(Some(worker)));
|
|
|
|
// Spawn the worker in the background before moving into state
|
|
// Take the worker out of the Mutex so we don't hold the lock during the loop
|
|
{
|
|
let mut w = worker_handle.lock().await;
|
|
if let Some(worker) = w.take() {
|
|
tokio::spawn(async move {
|
|
if let Err(e) = worker.start().await {
|
|
tracing::error!("Worker error: {}", e);
|
|
}
|
|
});
|
|
}
|
|
}
|
|
|
|
let state = AppState {
|
|
config: config.clone(),
|
|
queue,
|
|
worker: worker_handle,
|
|
};
|
|
|
|
let app = create_router(state);
|
|
|
|
let addr = config.server.listen.clone();
|
|
info!("Starting API server on {}", addr);
|
|
|
|
let listener = tokio::net::TcpListener::bind(&addr)
|
|
.await
|
|
.map_err(|e| crate::error::FaxError::Other(format!("Failed to bind: {}", e)))?;
|
|
|
|
axum::serve(listener, app)
|
|
.await
|
|
.map_err(|e| crate::error::FaxError::Other(format!("Server error: {}", e)))?;
|
|
|
|
Ok(())
|
|
} |