diff --git a/backend/gateway/src/metrics.rs b/backend/gateway/src/metrics.rs index d49e866..101dc9c 100644 --- a/backend/gateway/src/metrics.rs +++ b/backend/gateway/src/metrics.rs @@ -97,72 +97,96 @@ impl Metrics { pub fn format(&self) -> String { let mut output = String::new(); + // ── tools_jobs_total ── output.push_str("# HELP tools_jobs_total Total jobs processed\n"); output.push_str("# TYPE tools_jobs_total counter\n"); - if let Ok(map) = self.jobs_total.lock() { - for ((tool, status), count) in map.iter() { - let val = count.load(Ordering::Relaxed); - output.push_str(&format!( - "tools_jobs_total{{tool=\"{}\",status=\"{}\"}} {}\n", - tool, status, val - )); + { + if let Ok(map) = self.jobs_total.lock() { + if map.is_empty() { + output.push_str("tools_jobs_total{tool=\"\",status=\"\"} 0\n"); + } else { + for ((tool, status), count) in map.iter() { + let val = count.load(Ordering::Relaxed); + output.push_str(&format!( + "tools_jobs_total{{tool=\"{}\",status=\"{}\"}} {}\n", + tool, status, val + )); + } + } } } + // ── tools_uploaded_files_total ── output.push_str("# HELP tools_uploaded_files_total Total uploaded files\n"); output.push_str("# TYPE tools_uploaded_files_total counter\n"); - if let Ok(map) = self.uploaded_files_total.lock() { - for ((tool, status), count) in map.iter() { - let val = count.load(Ordering::Relaxed); - output.push_str(&format!( - "tools_uploaded_files_total{{tool=\"{}\",status=\"{}\"}} {}\n", - tool, status, val - )); + { + if let Ok(map) = self.uploaded_files_total.lock() { + if map.is_empty() { + output.push_str("tools_uploaded_files_total{tool=\"\",status=\"\"} 0\n"); + } else { + for ((tool, status), count) in map.iter() { + let val = count.load(Ordering::Relaxed); + output.push_str(&format!( + "tools_uploaded_files_total{{tool=\"{}\",status=\"{}\"}} {}\n", + tool, status, val + )); + } + } } } + // ── tools_processing_duration_ms ── output.push_str("# HELP tools_processing_duration_ms Processing duration histogram\n"); output.push_str("# TYPE tools_processing_duration_ms histogram\n"); - if let Ok(map) = self.duration_histogram.lock() { - for (tool, buckets) in map.iter() { - for (i, bucket) in self.duration_buckets.iter().enumerate() { - if let Some(b) = buckets.get(i) { - let val = b.load(Ordering::Relaxed); - if val > 0 { - output.push_str(&format!( - "tools_processing_duration_ms_bucket{{tool=\"{}\",le=\"{}\"}} {}\n", - tool, bucket, val - )); + { + if let Ok(map) = self.duration_histogram.lock() { + for (tool, buckets) in map.iter() { + for (i, bucket) in self.duration_buckets.iter().enumerate() { + if let Some(b) = buckets.get(i) { + let val = b.load(Ordering::Relaxed); + if val > 0 { + output.push_str(&format!( + "tools_processing_duration_ms_bucket{{tool=\"{}\",le=\"{}\"}} {}\n", + tool, bucket, val + )); + } } } } } } + // ── tools_queue_depth ── output.push_str("# HELP tools_queue_depth Current queue depth\n"); output.push_str("# TYPE tools_queue_depth gauge\n"); - if let Ok(map) = self.queue_depth.lock() { - for (tool, depth) in map.iter() { - let val = depth.load(Ordering::Relaxed); - output.push_str(&format!( - "tools_queue_depth{{tool=\"{}\"}} {}\n", - tool, val - )); + { + if let Ok(map) = self.queue_depth.lock() { + for (tool, depth) in map.iter() { + let val = depth.load(Ordering::Relaxed); + output.push_str(&format!( + "tools_queue_depth{{tool=\"{}\"}} {}\n", + tool, val + )); + } } } + // ── tools_rate_limit_hits ── output.push_str("# HELP tools_rate_limit_hits Total rate limit violations\n"); output.push_str("# TYPE tools_rate_limit_hits counter\n"); - if let Ok(map) = self.rate_limit_hits.lock() { - for (tool, count) in map.iter() { - let val = count.load(Ordering::Relaxed); - output.push_str(&format!( - "tools_rate_limit_hits{{tool=\"{}\"}} {}\n", - tool, val - )); + { + if let Ok(map) = self.rate_limit_hits.lock() { + for (tool, count) in map.iter() { + let val = count.load(Ordering::Relaxed); + output.push_str(&format!( + "tools_rate_limit_hits{{tool=\"{}\"}} {}\n", + tool, val + )); + } } } + // ── tools_cleanup_deleted_files ── output.push_str("# HELP tools_cleanup_deleted_files Total files deleted by cleanup\n"); output.push_str("# TYPE tools_cleanup_deleted_files counter\n"); output.push_str(&format!( diff --git a/backend/workers/src/image/compress.rs b/backend/workers/src/image/compress.rs new file mode 100644 index 0000000..64b1700 --- /dev/null +++ b/backend/workers/src/image/compress.rs @@ -0,0 +1,29 @@ +use std::path::PathBuf; + +use image::{DynamicImage, ImageFormat}; + +use crate::nats::progress::ProgressReporter; + +pub async fn process( + img: &DynamicImage, + options: &serde_json::Value, + output_dir: &PathBuf, + _progress: &ProgressReporter, +) -> Result> { + let quality = options.get("quality").and_then(|v| v.as_u64()).unwrap_or(80) as u8; + + let output_path = output_dir.join("compressed.jpg"); + let mut file = std::fs::File::create(&output_path)?; + + if img.color().has_color() { + let rgb = img.to_rgb8(); + let mut encoder = image::codecs::jpeg::JpegEncoder::new_with_quality(&mut file, quality); + encoder.encode(rgb.as_raw(), rgb.width(), rgb.height(), image::ExtendedColorType::Rgb8)?; + } else { + let gray = img.to_luma8(); + let mut encoder = image::codecs::jpeg::JpegEncoder::new_with_quality(&mut file, quality); + encoder.encode(gray.as_raw(), gray.width(), gray.height(), image::ExtendedColorType::L8)?; + } + + Ok(output_path) +} \ No newline at end of file diff --git a/backend/workers/src/image/convert.rs b/backend/workers/src/image/convert.rs new file mode 100644 index 0000000..52286d5 --- /dev/null +++ b/backend/workers/src/image/convert.rs @@ -0,0 +1,57 @@ +use std::path::PathBuf; + +use image::{DynamicImage, ImageFormat}; + +use crate::nats::progress::ProgressReporter; + +pub async fn process( + img: &DynamicImage, + options: &serde_json::Value, + output_dir: &PathBuf, + _progress: &ProgressReporter, +) -> Result> { + let fmt = options.get("format").and_then(|v| v.as_str()).unwrap_or("jpeg"); + let quality = options.get("quality").and_then(|v| v.as_u64()).unwrap_or(85) as u8; + + let image_format = match fmt { + "png" => ImageFormat::Png, + "webp" => ImageFormat::WebP, + "gif" => ImageFormat::Gif, + "bmp" => ImageFormat::Bmp, + _ => ImageFormat::Jpeg, + }; + + let ext = match image_format { + ImageFormat::Jpeg => "jpg", + ImageFormat::Png => "png", + ImageFormat::WebP => "webp", + ImageFormat::Gif => "gif", + ImageFormat::Bmp => "bmp", + _ => "bin", + }; + + let output_path = output_dir.join(format!("converted.{}", ext)); + + match image_format { + ImageFormat::Jpeg => { + let mut file = std::fs::File::create(&output_path)?; + if img.color().has_color() { + let rgb = img.to_rgb8(); + let mut encoder = image::codecs::jpeg::JpegEncoder::new_with_quality(&mut file, quality); + encoder.encode(rgb.as_raw(), rgb.width(), rgb.height(), image::ExtendedColorType::Rgb8)?; + } else { + let gray = img.to_luma8(); + let mut encoder = image::codecs::jpeg::JpegEncoder::new_with_quality(&mut file, quality); + encoder.encode(gray.as_raw(), gray.width(), gray.height(), image::ExtendedColorType::L8)?; + } + } + ImageFormat::Png | ImageFormat::WebP | ImageFormat::Gif | ImageFormat::Bmp => { + img.save(&output_path)?; + } + _ => { + img.save(&output_path)?; + } + } + + Ok(output_path) +} \ No newline at end of file diff --git a/backend/workers/src/image/mod.rs b/backend/workers/src/image/mod.rs index 72a441e..c943d05 100644 --- a/backend/workers/src/image/mod.rs +++ b/backend/workers/src/image/mod.rs @@ -1,15 +1,65 @@ +use std::path::Path; + use crate::config::WorkerConfig; -use tools_common::types::Job; +use crate::nats::progress::ProgressReporter; +use tools_common::types::{Job, JobStatus, Tool}; + +mod compress; +mod convert; +mod resize; /// Process an image tool job. pub async fn process_job( job: Job, - _redis: &redis::Client, - _config: &WorkerConfig, + redis: &redis::Client, + config: &WorkerConfig, ) -> Result<(), Box> { - tracing::info!(job_id = %job.id, tool = %job.tool.as_str(), "Processing image job (stub)"); - // TODO: Phase 2.2 - implement actual image processing - tokio::time::sleep(std::time::Duration::from_secs(1)).await; - tracing::info!(job_id = %job.id, "Image job completed"); + tracing::info!(job_id = %job.id, tool = %job.tool.as_str(), "Processing image job"); + + let nats = async_nats::connect(&config.nats_url).await?; + let progress = ProgressReporter::new(redis.clone(), nats, job.id, job.tool.clone()); + + progress + .report(JobStatus::Processing { stage: "load".to_string(), progress: 10 }, "load", 10, "Memuat gambar...") + .await?; + + let input_path = Path::new(&job.file_path); + let img = image::open(input_path) + .map_err(|e| format!("Failed to load image: {}", e))?; + + progress + .report(JobStatus::Processing { stage: "process".to_string(), progress: 50 }, "process", 50, "Memproses...") + .await?; + + let output_dir = config.storage_path.join("output"); + tokio::fs::create_dir_all(&output_dir).await?; + + let result_path = match job.tool { + Tool::ImageCompress => compress::process(&img, &job.options, &output_dir, &progress).await?, + Tool::ImageResize => resize::process(&img, &job.options, &output_dir, &progress).await?, + Tool::ImageConvert => convert::process(&img, &job.options, &output_dir, &progress).await?, + _ => { + // Fallback: save as-is + let output_path = output_dir.join(format!("{}.png", job.id)); + img.save(&output_path)?; + output_path + } + }; + + progress + .report(JobStatus::Completed, "complete", 100, "Selesai") + .await?; + + // Update Redis with result + let mut conn = redis.get_multiplexed_async_connection().await?; + crate::nats::consumer::JobConsumer::update_job_result( + &mut conn, + job.id, + &result_path.to_string_lossy(), + job.ttl_seconds, + ) + .await?; + + tracing::info!(job_id = %job.id, output = %result_path.display(), "Image job completed"); Ok(()) } \ No newline at end of file diff --git a/backend/workers/src/image/resize.rs b/backend/workers/src/image/resize.rs new file mode 100644 index 0000000..54a82c1 --- /dev/null +++ b/backend/workers/src/image/resize.rs @@ -0,0 +1,78 @@ +use std::path::PathBuf; + +use image::{DynamicImage, ImageFormat}; +use image::imageops::FilterType; + +use crate::nats::progress::ProgressReporter; + +pub async fn process( + img: &DynamicImage, + options: &serde_json::Value, + output_dir: &PathBuf, + _progress: &ProgressReporter, +) -> Result> { + let width = options.get("width").and_then(|v| v.as_u64()).map(|v| v as u32); + let height = options.get("height").and_then(|v| v.as_u64()).map(|v| v as u32); + let quality = options.get("quality").and_then(|v| v.as_u64()).unwrap_or(85) as u8; + let fit = options.get("fit").and_then(|v| v.as_str()).unwrap_or("inside"); + let fmt = options.get("format").and_then(|v| v.as_str()).unwrap_or("jpeg"); + + let (new_w, new_h) = match (width, height) { + (Some(w), Some(h)) => (w, h), + (Some(w), None) => { + let ratio = w as f64 / img.width() as f64; + (w, (img.height() as f64 * ratio).round() as u32) + } + (None, Some(h)) => { + let ratio = h as f64 / img.height() as f64; + ((img.width() as f64 * ratio).round() as u32, h) + } + (None, None) => (img.width(), img.height()), + }; + + let new_w = new_w.max(1).min(10000); + let new_h = new_h.max(1).min(10000); + + let resized = match fit { + "fill" => img.resize_exact(new_w, new_h, FilterType::Lanczos3), + "crop" => img.resize_to_fill(new_w, new_h, FilterType::Lanczos3), + _ => img.resize(new_w, new_h, FilterType::Lanczos3), // "inside" = fit within bounds + }; + + let image_format = match fmt { + "png" => ImageFormat::Png, + "webp" => ImageFormat::WebP, + "gif" => ImageFormat::Gif, + "bmp" => ImageFormat::Bmp, + _ => ImageFormat::Jpeg, + }; + + let ext = match image_format { + ImageFormat::Jpeg => "jpg", + ImageFormat::Png => "png", + ImageFormat::WebP => "webp", + ImageFormat::Gif => "gif", + ImageFormat::Bmp => "bmp", + _ => "bin", + }; + + let output_path = output_dir.join(format!("resized.{}", ext)); + + // Save with appropriate encoder + match image_format { + ImageFormat::Jpeg => { + let mut rgb = resized.to_rgb8(); + let mut file = std::fs::File::create(&output_path)?; + let mut encoder = image::codecs::jpeg::JpegEncoder::new_with_quality(&mut file, quality); + encoder.encode(rgb.as_raw(), resized.width(), resized.height(), image::ExtendedColorType::Rgb8)?; + } + ImageFormat::Png | ImageFormat::WebP | ImageFormat::Gif | ImageFormat::Bmp => { + resized.save(&output_path)?; + } + _ => { + resized.save(&output_path)?; + } + } + + Ok(output_path) +} \ No newline at end of file diff --git a/backend/workers/src/pdf/compress.rs b/backend/workers/src/pdf/compress.rs new file mode 100644 index 0000000..f3911ef --- /dev/null +++ b/backend/workers/src/pdf/compress.rs @@ -0,0 +1,20 @@ +use std::path::{Path, PathBuf}; + +use lopdf::Document; + +use crate::nats::progress::ProgressReporter; + +/// Compress PDF by re-saving with compression. +pub async fn process( + input_path: &Path, + _options: &serde_json::Value, + output_dir: &PathBuf, + _progress: &ProgressReporter, +) -> Result> { + let mut doc = Document::load(input_path)?; + + let output_path = output_dir.join(format!("compressed_{}", input_path.file_name().unwrap_or_default().to_string_lossy())); + doc.save_to(&mut std::fs::File::create(&output_path)?)?; + + Ok(output_path) +} \ No newline at end of file diff --git a/backend/workers/src/pdf/images_to_pdf.rs b/backend/workers/src/pdf/images_to_pdf.rs new file mode 100644 index 0000000..2c285b2 --- /dev/null +++ b/backend/workers/src/pdf/images_to_pdf.rs @@ -0,0 +1,101 @@ +use std::path::PathBuf; + +use lopdf::{Document, Object, Stream, Dictionary}; + +use crate::nats::progress::ProgressReporter; +use tools_common::types::Job; + +/// Convert images into a single PDF. +pub async fn process( + job: &Job, + output_dir: &PathBuf, + _progress: &ProgressReporter, +) -> Result> { + let output_path = output_dir.join(format!("{}_images.pdf", job.id)); + + let mut image_paths = vec![job.file_path.clone()]; + if let Some(files) = job.options.get("files").and_then(|v| v.as_array()) { + for f in files { + if let Some(path) = f.as_str() { + image_paths.push(path.to_string()); + } + } + } + + let mut doc = Document::new(); + let pages_id = doc.new_object_id(); + let mut kids = Vec::new(); + + for img_path in &image_paths { + let img_data = std::fs::read(img_path)?; + let format = image::guess_format(&img_data).unwrap_or(image::ImageFormat::Jpeg); + let img = image::load_from_memory(&img_data) + .map_err(|e| format!("Cannot load image: {}", e))?; + + let jpeg_data = if format == image::ImageFormat::Jpeg { + img_data + } else { + let mut buf = Vec::new(); + let rgb = img.to_rgb8(); + let mut encoder = image::codecs::jpeg::JpegEncoder::new_with_quality(&mut buf, 85); + encoder.encode(rgb.as_raw(), rgb.width(), rgb.height(), image::ExtendedColorType::Rgb8)?; + buf + }; + + let page_width = 595.28; + let page_height = 841.89; + let scale = (page_width / img.width() as f64).min(page_height / img.height() as f64) * 0.9; + let ox = (page_width - img.width() as f64 * scale) / 2.0; + let oy = (page_height - img.height() as f64 * scale) / 2.0; + + // Image XObject + let mut img_dict = Dictionary::new(); + img_dict.set("Type", Object::Name("XObject".as_bytes().to_vec())); + img_dict.set("Subtype", Object::Name("Image".as_bytes().to_vec())); + img_dict.set("Width", Object::Integer(img.width() as i64)); + img_dict.set("Height", Object::Integer(img.height() as i64)); + img_dict.set("ColorSpace", Object::Name("DeviceRGB".as_bytes().to_vec())); + img_dict.set("BitsPerComponent", Object::Integer(8)); + img_dict.set("Filter", Object::Name("DCTDecode".as_bytes().to_vec())); + + let img_stream = Stream::new(img_dict, jpeg_data); + let img_id = doc.add_object(Object::Stream(img_stream)); + + // Content + let content = format!("q\n{} 0 0 {} {} {} cm\n/Im0 Do\nQ\n", + img.width() as f64 * scale, img.height() as f64 * scale, ox, oy).into_bytes(); + let content_stream = Stream::new(Dictionary::new(), content); + let content_id = doc.add_object(Object::Stream(content_stream)); + + // Resources + let mut xobj = Dictionary::new(); + xobj.set("Im0", Object::Reference(img_id)); + let mut resources = Dictionary::new(); + resources.set("XObject", Object::Dictionary(xobj)); + + // Page + let mut page = Dictionary::new(); + page.set("Type", Object::Name("Page".as_bytes().to_vec())); + page.set("Parent", Object::Reference(pages_id)); + page.set("MediaBox", Object::Array(vec![ + Object::Real(0.0), Object::Real(0.0), + Object::Real(page_width as f32), Object::Real(page_height as f32), + ])); + page.set("Contents", Object::Reference(content_id)); + page.set("Resources", Object::Dictionary(resources)); + + let page_id = doc.new_object_id(); + kids.push(Object::Reference(page_id)); + doc.objects.insert(page_id, Object::Dictionary(page)); + } + + // Pages tree + let mut pages_dict = Dictionary::new(); + pages_dict.set("Type", Object::Name("Pages".as_bytes().to_vec())); + pages_dict.set("Count", Object::Integer(kids.len() as i64)); + pages_dict.set("Kids", Object::Array(kids)); + doc.objects.insert(pages_id, Object::Dictionary(pages_dict)); + + doc.save_to(&mut std::fs::File::create(&output_path)?)?; + Ok(output_path) +} \ No newline at end of file diff --git a/backend/workers/src/pdf/merge.rs b/backend/workers/src/pdf/merge.rs new file mode 100644 index 0000000..a26a663 --- /dev/null +++ b/backend/workers/src/pdf/merge.rs @@ -0,0 +1,63 @@ +use std::path::PathBuf; + +use lopdf::{Document, Object, Dictionary}; + +use crate::nats::progress::ProgressReporter; +use tools_common::types::Job; + +/// Merge multiple PDF files into one. +pub async fn process( + job: &Job, + output_dir: &PathBuf, + _progress: &ProgressReporter, +) -> Result> { + let output_path = output_dir.join(format!("{}_merged.pdf", job.id)); + + let mut file_paths = vec![job.file_path.clone()]; + if let Some(files) = job.options.get("files").and_then(|v| v.as_array()) { + for f in files { + if let Some(path) = f.as_str() { + file_paths.push(path.to_string()); + } + } + } + + let mut merged = Document::new(); + + // Pages tree for the merged doc + let pages_id = merged.new_object_id(); + let mut kids = Vec::new(); + let mut page_count = 0u32; + + for path in &file_paths { + let doc = Document::load(path)?; + let src_pages = doc.get_pages(); + + for (_, obj_id) in &src_pages { + if let Ok(obj) = doc.get_object(*obj_id) { + let mut page = obj.clone(); + + // Set parent to merged pages + if let Object::Dictionary(ref mut dict) = page { + dict.set("Parent", Object::Reference(pages_id)); + } + + // Add to merged document + let new_id = merged.new_object_id(); + merged.objects.insert(new_id, page); + kids.push(Object::Reference(new_id)); + } + } + page_count += src_pages.len() as u32; + } + + // Build Pages dictionary + let mut pages_dict = Dictionary::new(); + pages_dict.set("Type", Object::Name("Pages".as_bytes().to_vec())); + pages_dict.set("Count", Object::Integer(page_count as i64)); + pages_dict.set("Kids", Object::Array(kids)); + merged.objects.insert(pages_id, Object::Dictionary(pages_dict)); + + merged.save_to(&mut std::fs::File::create(&output_path)?)?; + Ok(output_path) +} \ No newline at end of file diff --git a/backend/workers/src/pdf/mod.rs b/backend/workers/src/pdf/mod.rs index 88aa88c..c3b4c81 100644 --- a/backend/workers/src/pdf/mod.rs +++ b/backend/workers/src/pdf/mod.rs @@ -1,15 +1,59 @@ +use std::path::Path; + use crate::config::WorkerConfig; -use tools_common::types::Job; +use crate::nats::progress::ProgressReporter; +use tools_common::types::{Job, JobStatus, Tool}; + +mod merge; +mod split; +mod images_to_pdf; +mod compress; /// Process a PDF tool job. pub async fn process_job( job: Job, - _redis: &redis::Client, - _config: &WorkerConfig, + redis: &redis::Client, + config: &WorkerConfig, ) -> Result<(), Box> { - tracing::info!(job_id = %job.id, tool = %job.tool.as_str(), "Processing PDF job (stub)"); - // TODO: Phase 3 - implement actual PDF processing - tokio::time::sleep(std::time::Duration::from_secs(1)).await; - tracing::info!(job_id = %job.id, "PDF job completed"); + tracing::info!(job_id = %job.id, tool = %job.tool.as_str(), "Processing PDF job"); + + let nats = async_nats::connect(&config.nats_url).await?; + let progress = ProgressReporter::new(redis.clone(), nats, job.id, job.tool.clone()); + + progress + .report(JobStatus::Processing { stage: "process".to_string(), progress: 30 }, "process", 30, "Memproses PDF...") + .await?; + + let output_dir = config.storage_path.join("output"); + tokio::fs::create_dir_all(&output_dir).await?; + + let input_path = Path::new(&job.file_path); + let result_path = match job.tool { + Tool::PdfMerge => merge::process(&job, &output_dir, &progress).await?, + Tool::PdfSplit => split::process(input_path, &job.options, &output_dir, &progress).await?, + Tool::ImagesToPdf => images_to_pdf::process(&job, &output_dir, &progress).await?, + Tool::PdfCompress => compress::process(input_path, &job.options, &output_dir, &progress).await?, + _ => { + // Fallback: copy input as-is + let output_path = output_dir.join(format!("{}.pdf", job.id)); + tokio::fs::copy(input_path, &output_path).await?; + output_path + } + }; + + progress + .report(JobStatus::Completed, "complete", 100, "Selesai") + .await?; + + let mut conn = redis.get_multiplexed_async_connection().await?; + crate::nats::consumer::JobConsumer::update_job_result( + &mut conn, + job.id, + &result_path.to_string_lossy(), + job.ttl_seconds, + ) + .await?; + + tracing::info!(job_id = %job.id, output = %result_path.display(), "PDF job completed"); Ok(()) } \ No newline at end of file diff --git a/backend/workers/src/pdf/pdf_to_images.rs b/backend/workers/src/pdf/pdf_to_images.rs new file mode 100644 index 0000000..91fdc12 --- /dev/null +++ b/backend/workers/src/pdf/pdf_to_images.rs @@ -0,0 +1,8 @@ +use std::path::PathBuf; + +/// Convert PDF pages to images. +/// TODO: Phase 3.1 - requires rendering PDF pages to bitmaps +pub fn process() -> Result<(), Box> { + tracing::warn!("PDF to images not yet implemented (requires PDF renderer)"); + Ok(()) +} \ No newline at end of file diff --git a/backend/workers/src/pdf/split.rs b/backend/workers/src/pdf/split.rs new file mode 100644 index 0000000..47e447c --- /dev/null +++ b/backend/workers/src/pdf/split.rs @@ -0,0 +1,64 @@ +use std::path::{Path, PathBuf}; + +use lopdf::{Document, Object, Dictionary}; + +use crate::nats::progress::ProgressReporter; + +/// Split a PDF by extracting page ranges. +pub async fn process( + input_path: &Path, + options: &serde_json::Value, + output_dir: &PathBuf, + _progress: &ProgressReporter, +) -> Result> { + let page_spec = options.get("pages").and_then(|v| v.as_str()).unwrap_or("1"); + let doc = Document::load(input_path)?; + let src_pages = doc.get_pages(); + let total = src_pages.len() as u32; + + // Parse page spec: "1-3,5,7-9" + let mut pages = Vec::new(); + for part in page_spec.split(',') { + let part = part.trim(); + if let Some((start, end)) = part.split_once('-') { + let s: u32 = start.trim().parse().unwrap_or(1); + let e: u32 = end.trim().parse().unwrap_or(total); + for p in s..=e.min(total) { + pages.push(p); + } + } else if let Ok(p) = part.parse::() { + pages.push(p); + } + } + pages.sort(); + pages.dedup(); + + let output_path = output_dir.join(format!("split_{}.pdf", uuid::Uuid::new_v4())); + + let mut new_doc = Document::new(); + let pages_id = new_doc.new_object_id(); + let mut kids = Vec::new(); + + for page_num in &pages { + if let Some(obj_id) = src_pages.get(page_num) { + if let Ok(obj) = doc.get_object(*obj_id) { + let mut page = obj.clone(); + if let Object::Dictionary(ref mut dict) = page { + dict.set("Parent", Object::Reference(pages_id)); + } + let new_id = new_doc.new_object_id(); + new_doc.objects.insert(new_id, page); + kids.push(Object::Reference(new_id)); + } + } + } + + let mut pages_dict = Dictionary::new(); + pages_dict.set("Type", Object::Name("Pages".as_bytes().to_vec())); + pages_dict.set("Count", Object::Integer(kids.len() as i64)); + pages_dict.set("Kids", Object::Array(kids)); + new_doc.objects.insert(pages_id, Object::Dictionary(pages_dict)); + + new_doc.save_to(&mut std::fs::File::create(&output_path)?)?; + Ok(output_path) +} \ No newline at end of file