main.rs (3663B)
1 #![warn(clippy::missing_const_for_fn)] 2 #![warn(clippy::suboptimal_flops)] 3 4 use std::path::PathBuf; 5 6 use async_nats::jetstream::{self, stream}; 7 use async_nats::jetstream::stream::Compression; 8 use clap::{Parser, Subcommand}; 9 use notify_debouncer_full::notify::{RecursiveMode, Watcher}; 10 use tracing::info; 11 12 mod c2pa; 13 mod importer; 14 mod jetstream_events; 15 16 #[derive(clap::Parser)] 17 struct Opts { 18 #[arg(short, long, default_value = "nats://localhost:4222")] 19 nats_url: String, 20 #[arg(short, long, default_value = "import")] 21 import_folder: PathBuf, 22 #[arg(short, long, default_value = "export")] 23 export_folder: PathBuf, 24 #[arg(short, long, default_value = "c2pa")] 25 c2pa_folder: PathBuf, 26 #[arg(long, default_value = "true")] 27 walk_export_folders: std::primitive::bool, 28 #[arg(long, default_value = "true")] 29 walk_import_folders: std::primitive::bool, 30 31 // Subcommands 32 #[command(subcommand)] 33 command: Option<Commands>, 34 } 35 36 #[derive(Subcommand)] 37 enum Commands { 38 C2PACert {}, 39 } 40 41 #[tokio::main] 42 async fn main() -> color_eyre::Result<()> { 43 color_eyre::install()?; 44 45 // Set default login to default settings (info, warn, error being shown) 46 tracing_subscriber::fmt::Subscriber::builder() 47 .with_env_filter("info") 48 .init(); 49 50 // Parse command line arguments 51 let opts: Opts = Opts::parse(); 52 53 // Create import folder if it doesn't exist 54 let import_folder = opts.import_folder; 55 let export_folder = opts.export_folder; 56 let c2pa_folder = opts.c2pa_folder; 57 58 std::fs::create_dir_all(import_folder.clone())?; 59 std::fs::create_dir_all(export_folder.clone())?; 60 std::fs::create_dir_all(c2pa_folder.clone())?; 61 62 if let Some(Commands::C2PACert {}) = opts.command { 63 c2pa::generate_certificates(c2pa_folder.clone())?; 64 return Ok(()); 65 } 66 67 rexiv2::initialize().expect("Unable to initialize rexiv2"); 68 69 info!("Starting the application"); 70 71 // Setup nats connection 72 let nats_url = std::env::var("NATS_URL").unwrap_or(opts.nats_url); 73 74 let client = async_nats::connect(nats_url).await?; 75 76 let jetstream = jetstream::new(client); 77 78 let import_stream = jetstream 79 .create_stream(stream::Config { 80 name: "IMPORT".to_string(), 81 retention: stream::RetentionPolicy::WorkQueue, 82 subjects: vec!["import.>".to_string()], 83 compression: Some(Compression::S2), 84 ..Default::default() 85 }) 86 .await?; 87 info!("Nats import stream created"); 88 89 info!("Creating C2PA signer"); 90 let c2pa_signer = 91 c2pa::C2PASigner::new(export_folder.clone(), c2pa_folder, import_stream.clone()).await?; 92 tokio::spawn(async move { 93 c2pa_signer.consumer().await.unwrap(); 94 }); 95 96 info!("Creating importer in other thread"); 97 98 let importer = importer::Importer::new( 99 import_folder.clone(), 100 export_folder, 101 import_stream, 102 jetstream, 103 ) 104 .await?; 105 let importer_clone = importer.clone(); 106 tokio::spawn(async move { 107 importer_clone.nats_consumer().await.unwrap(); 108 }); 109 110 // Create a file listener using notify to listen for new files to import 111 let mut debouncer = importer 112 .listen_for_files(opts.walk_import_folders, opts.walk_export_folders) 113 .await 114 .unwrap(); 115 debouncer 116 .watcher() 117 .watch(&import_folder, RecursiveMode::NonRecursive)?; 118 debouncer 119 .cache() 120 .add_root(&import_folder, RecursiveMode::NonRecursive); 121 info!("Watching for new files in {:?}", import_folder); 122 123 // Wait for ctrl-c to exit 124 tokio::signal::ctrl_c().await?; 125 126 Ok(()) 127 }