diff --git a/core/src/api/files.rs b/core/src/api/files.rs index 2fc1f27a7..61b4c9a04 100644 --- a/core/src/api/files.rs +++ b/core/src/api/files.rs @@ -1,7 +1,7 @@ use crate::{ invalidate_query, job::Job, - library::LibraryContext, + library::Library, object::fs::{ copy::{FileCopierJob, FileCopierJobInit}, cut::{FileCutterJob, FileCutterJobInit}, @@ -26,7 +26,7 @@ pub(crate) fn mount() -> RouterBuilder { pub struct GetArgs { pub id: i32, } - t(|_, args: GetArgs, library: LibraryContext| async move { + t(|_, args: GetArgs, library: Library| async move { Ok(library .db .object() @@ -43,7 +43,7 @@ pub(crate) fn mount() -> RouterBuilder { pub note: Option, } - t(|_, args: SetNoteArgs, library: LibraryContext| async move { + t(|_, args: SetNoteArgs, library: Library| async move { library .db .object() @@ -67,27 +67,25 @@ pub(crate) fn mount() -> RouterBuilder { pub favorite: bool, } - t( - |_, args: SetFavoriteArgs, library: LibraryContext| async move { - library - .db - .object() - .update( - object::id::equals(args.id), - vec![object::favorite::set(args.favorite)], - ) - .exec() - .await?; + t(|_, args: SetFavoriteArgs, library: Library| async move { + library + .db + .object() + .update( + object::id::equals(args.id), + vec![object::favorite::set(args.favorite)], + ) + .exec() + .await?; - invalidate_query!(library, "locations.getExplorerData"); - invalidate_query!(library, "tags.getExplorerData"); + invalidate_query!(library, "locations.getExplorerData"); + invalidate_query!(library, "tags.getExplorerData"); - Ok(()) - }, - ) + Ok(()) + }) }) .library_mutation("delete", |t| { - t(|_, id: i32, library: LibraryContext| async move { + t(|_, id: i32, library: Library| async move { library .db .object() @@ -101,7 +99,7 @@ pub(crate) fn mount() -> RouterBuilder { }) .library_mutation("encryptFiles", |t| { t( - |_, args: FileEncryptorJobInit, library: LibraryContext| async move { + |_, args: FileEncryptorJobInit, library: Library| async move { library.spawn_job(Job::new(args, FileEncryptorJob {})).await; invalidate_query!(library, "locations.getExplorerData"); @@ -111,7 +109,7 @@ pub(crate) fn mount() -> RouterBuilder { }) .library_mutation("decryptFiles", |t| { t( - |_, args: FileDecryptorJobInit, library: LibraryContext| async move { + |_, args: FileDecryptorJobInit, library: Library| async move { library.spawn_job(Job::new(args, FileDecryptorJob {})).await; invalidate_query!(library, "locations.getExplorerData"); @@ -120,75 +118,65 @@ pub(crate) fn mount() -> RouterBuilder { ) }) .library_mutation("deleteFiles", |t| { - t( - |_, args: FileDeleterJobInit, library: LibraryContext| async move { - library.spawn_job(Job::new(args, FileDeleterJob {})).await; - invalidate_query!(library, "locations.getExplorerData"); + t(|_, args: FileDeleterJobInit, library: Library| async move { + library.spawn_job(Job::new(args, FileDeleterJob {})).await; + invalidate_query!(library, "locations.getExplorerData"); - Ok(()) - }, - ) + Ok(()) + }) }) .library_mutation("eraseFiles", |t| { - t( - |_, args: FileEraserJobInit, library: LibraryContext| async move { - library.spawn_job(Job::new(args, FileEraserJob {})).await; - invalidate_query!(library, "locations.getExplorerData"); + t(|_, args: FileEraserJobInit, library: Library| async move { + library.spawn_job(Job::new(args, FileEraserJob {})).await; + invalidate_query!(library, "locations.getExplorerData"); - Ok(()) - }, - ) + Ok(()) + }) }) .library_mutation("duplicateFiles", |t| { - t( - |_, args: FileCopierJobInit, library: LibraryContext| async move { - let (done_tx, done_rx) = oneshot::channel(); + t(|_, args: FileCopierJobInit, library: Library| async move { + let (done_tx, done_rx) = oneshot::channel(); - library - .spawn_job(Job::new( - args, - FileCopierJob { - done_tx: Some(done_tx), - }, - )) - .await; + library + .spawn_job(Job::new( + args, + FileCopierJob { + done_tx: Some(done_tx), + }, + )) + .await; - let _ = done_rx.await; - invalidate_query!(library, "locations.getExplorerData"); + let _ = done_rx.await; + invalidate_query!(library, "locations.getExplorerData"); - Ok(()) - }, - ) + Ok(()) + }) }) .library_mutation("copyFiles", |t| { - t( - |_, args: FileCopierJobInit, library: LibraryContext| async move { - let (done_tx, done_rx) = oneshot::channel(); + t(|_, args: FileCopierJobInit, library: Library| async move { + let (done_tx, done_rx) = oneshot::channel(); - library - .spawn_job(Job::new( - args, - FileCopierJob { - done_tx: Some(done_tx), - }, - )) - .await; + library + .spawn_job(Job::new( + args, + FileCopierJob { + done_tx: Some(done_tx), + }, + )) + .await; - let _ = done_rx.await; - invalidate_query!(library, "locations.getExplorerData"); + let _ = done_rx.await; + invalidate_query!(library, "locations.getExplorerData"); - Ok(()) - }, - ) + Ok(()) + }) }) .library_mutation("cutFiles", |t| { - t( - |_, args: FileCutterJobInit, library: LibraryContext| async move { - library.spawn_job(Job::new(args, FileCutterJob {})).await; - invalidate_query!(library, "locations.getExplorerData"); + t(|_, args: FileCutterJobInit, library: Library| async move { + library.spawn_job(Job::new(args, FileCutterJob {})).await; + invalidate_query!(library, "locations.getExplorerData"); - Ok(()) - }, - ) + Ok(()) + }) }) } diff --git a/core/src/api/libraries.rs b/core/src/api/libraries.rs index 2eb13ef51..1cd13d741 100644 --- a/core/src/api/libraries.rs +++ b/core/src/api/libraries.rs @@ -1,7 +1,7 @@ use crate::{ api::Ctx, invalidate_query, - library::{LibraryConfig, LibraryContext}, + library::{Library, LibraryConfig}, prisma::statistics, volume::{get_volumes, save_volume}, }; @@ -28,7 +28,7 @@ pub(crate) fn mount() -> RouterBuilder { t(|ctx: Ctx, _: ()| async move { ctx.library_manager.get_all_libraries_config().await }) }) .library_query("getStatistics", |t| { - t(|_, _: (), library: LibraryContext| async move { + t(|_, _: (), library: Library| async move { let _statistics = library .db .statistics() diff --git a/core/src/api/locations.rs b/core/src/api/locations.rs index 5038f76d1..f4ac5b467 100644 --- a/core/src/api/locations.rs +++ b/core/src/api/locations.rs @@ -1,5 +1,5 @@ use crate::{ - library::LibraryContext, + library::Library, location::{ delete_location, fetch_location, indexer::{indexer_job::indexer_job_location, rules::IndexerRuleCreateArgs}, @@ -80,7 +80,7 @@ pub(crate) fn mount() -> impl RouterBuilderLike { } t(|_, mut args: LocationExplorerArgs, library| async move { - let LibraryContext { db, .. } = &library; + let Library { db, .. } = &library; let location = db .location() diff --git a/core/src/api/tags.rs b/core/src/api/tags.rs index dfcf500b9..4e797af05 100644 --- a/core/src/api/tags.rs +++ b/core/src/api/tags.rs @@ -8,7 +8,7 @@ use uuid::Uuid; use crate::{ api::locations::{object_with_file_paths, ExplorerContext, ExplorerData, ExplorerItem}, invalidate_query, - library::LibraryContext, + library::Library, prisma::{object, tag, tag_on_object}, sync, }; @@ -26,7 +26,7 @@ pub(crate) fn mount() -> RouterBuilder { t(|_, tag_id: i32, library| async move { info!("Getting files for tag {}", tag_id); - let LibraryContext { db, .. } = &library; + let Library { db, .. } = &library; let tag = db .tag() @@ -140,7 +140,7 @@ pub(crate) fn mount() -> RouterBuilder { } t(|_, args: TagCreateArgs, library| async move { - let LibraryContext { db, sync, .. } = &library; + let Library { db, sync, .. } = &library; let pub_id = Uuid::new_v4().as_bytes().to_vec(); @@ -211,7 +211,7 @@ pub(crate) fn mount() -> RouterBuilder { } t(|_, args: TagUpdateArgs, library| async move { - let LibraryContext { sync, db, .. } = &library; + let Library { sync, db, .. } = &library; let tag = db .tag() diff --git a/core/src/api/utils/invalidate.rs b/core/src/api/utils/invalidate.rs index 93eceb757..7d4d63b48 100644 --- a/core/src/api/utils/invalidate.rs +++ b/core/src/api/utils/invalidate.rs @@ -91,8 +91,8 @@ impl InvalidRequests { #[macro_export] #[allow(clippy::crate_in_macro_def)] macro_rules! invalidate_query { - ($ctx:expr, $key:literal) => {{ - let ctx: &crate::library::LibraryContext = &$ctx; // Assert the context is the correct type + ($library:expr, $key:literal) => {{ + let library: &crate::library::Library = &$library; // Assert the library is the correct type #[cfg(debug_assertions)] { @@ -111,13 +111,13 @@ macro_rules! invalidate_query { } // The error are ignored here because they aren't mission critical. If they fail the UI might be outdated for a bit. - ctx.emit(crate::api::CoreEvent::InvalidateOperation( + library.emit(crate::api::CoreEvent::InvalidateOperation( crate::api::utils::InvalidateOperationEvent::dangerously_create($key, serde_json::Value::Null) )) }}; - ($ctx:expr, $key:literal: $input_ty:ty, $input:expr $(,)?) => {{ + ($library:expr, $key:literal: $input_ty:ty, $input:expr $(,)?) => {{ let _: $input_ty = $input; // Assert the type the user provided is correct - let ctx: &crate::library::LibraryContext = &$ctx; // Assert the context is the correct type + let library: &crate::library::Library = &$library; // Assert the library is the correct type #[cfg(debug_assertions)] { @@ -141,7 +141,7 @@ macro_rules! invalidate_query { // The error are ignored here because they aren't mission critical. If they fail the UI might be outdated for a bit. let _ = serde_json::to_value($input) .map(|v| - ctx.emit(crate::api::CoreEvent::InvalidateOperation( + library.emit(crate::api::CoreEvent::InvalidateOperation( crate::api::utils::InvalidateOperationEvent::dangerously_create($key, v), )) ) diff --git a/core/src/api/utils/library.rs b/core/src/api/utils/library.rs index 893ddc87d..4de962fa8 100644 --- a/core/src/api/utils/library.rs +++ b/core/src/api/utils/library.rs @@ -11,7 +11,7 @@ use rspc::{ use serde::{de::DeserializeOwned, Deserialize, Serialize}; use uuid::Uuid; -use crate::{api::Ctx, library::LibraryContext}; +use crate::{api::Ctx, library::Library}; /// Can wrap a query argument to require it to contain a `library_id` and provide helpers for working with libraries. #[derive(Clone, Serialize, Deserialize, Type)] @@ -30,8 +30,8 @@ pub trait LibraryRequest { ) -> BuiltProcedureBuilder, ) -> Self where - TUnbuiltResolver: Fn(Ctx, TArg, LibraryContext) -> TUnbuiltResult + Send, - TBuiltResolver: Fn(Ctx, TArg, LibraryContext) -> TUnbuiltResult + Send + Sync + 'static, + TUnbuiltResolver: Fn(Ctx, TArg, Library) -> TUnbuiltResult + Send, + TBuiltResolver: Fn(Ctx, TArg, Library) -> TUnbuiltResult + Send + Sync + 'static, TUnbuiltResult: RequestResult + Send, TArg: DeserializeOwned + specta::Type + Send + 'static; @@ -49,8 +49,8 @@ pub trait LibraryRequest { ) -> BuiltProcedureBuilder, ) -> Self where - TUnbuiltResolver: Fn(Ctx, TArg, LibraryContext) -> TUnbuiltResult + Send, - TBuiltResolver: Fn(Ctx, TArg, LibraryContext) -> TUnbuiltResult + Send + Sync + 'static, + TUnbuiltResolver: Fn(Ctx, TArg, Library) -> TUnbuiltResult + Send, + TBuiltResolver: Fn(Ctx, TArg, Library) -> TUnbuiltResult + Send + Sync + 'static, TUnbuiltResult: RequestResult + Send, TArg: DeserializeOwned + specta::Type + Send + 'static; @@ -79,8 +79,8 @@ where ) -> BuiltProcedureBuilder, ) -> Self where - TUnbuiltResolver: Fn(Ctx, TArg, LibraryContext) -> TUnbuiltResult + Send, - TBuiltResolver: Fn(Ctx, TArg, LibraryContext) -> TUnbuiltResult + Send + Sync + 'static, + TUnbuiltResolver: Fn(Ctx, TArg, Library) -> TUnbuiltResult + Send, + TBuiltResolver: Fn(Ctx, TArg, Library) -> TUnbuiltResult + Send + Sync + 'static, TUnbuiltResult: RequestResult + Send, TArg: DeserializeOwned + specta::Type + Send + 'static, { @@ -125,8 +125,8 @@ where ) -> BuiltProcedureBuilder, ) -> Self where - TUnbuiltResolver: Fn(Ctx, TArg, LibraryContext) -> TUnbuiltResult + Send, - TBuiltResolver: Fn(Ctx, TArg, LibraryContext) -> TUnbuiltResult + Send + Sync + 'static, + TUnbuiltResolver: Fn(Ctx, TArg, Library) -> TUnbuiltResult + Send, + TBuiltResolver: Fn(Ctx, TArg, Library) -> TUnbuiltResult + Send + Sync + 'static, TUnbuiltResult: RequestResult + Send, TArg: DeserializeOwned + specta::Type + Send + 'static, { diff --git a/core/src/job/job_manager.rs b/core/src/job/job_manager.rs index 4858d68c4..22a1f2bd3 100644 --- a/core/src/job/job_manager.rs +++ b/core/src/job/job_manager.rs @@ -1,7 +1,7 @@ use crate::{ invalidate_query, job::{worker::Worker, DynJob, Job, JobError}, - library::LibraryContext, + library::Library, location::indexer::indexer_job::{IndexerJob, INDEXER_JOB_NAME}, object::{ fs::{ @@ -40,7 +40,7 @@ use uuid::Uuid; const MAX_WORKERS: usize = 1; pub enum JobManagerEvent { - IngestJob(LibraryContext, Box), + IngestJob(Library, Box), } /// JobManager handles queueing and executing jobs using the `DynJob` @@ -83,7 +83,7 @@ impl JobManager { this } - pub async fn ingest(self: Arc, ctx: &LibraryContext, job: Box) { + pub async fn ingest(self: Arc, ctx: &Library, job: Box) { let job_hash = job.hash(); debug!( "Ingesting job: ", @@ -119,7 +119,7 @@ impl JobManager { } } - pub async fn complete(self: Arc, ctx: &LibraryContext, job_id: Uuid, job_hash: u64) { + pub async fn complete(self: Arc, ctx: &Library, job_id: Uuid, job_hash: u64) { // remove worker from running workers and from current jobs hashes self.current_jobs_hashes.write().await.remove(&job_hash); self.running_workers.write().await.remove(&job_id); @@ -146,7 +146,7 @@ impl JobManager { } pub async fn get_history( - ctx: &LibraryContext, + ctx: &Library, ) -> Result, prisma_client_rust::QueryError> { Ok(ctx .db @@ -161,9 +161,7 @@ impl JobManager { .collect()) } - pub async fn clear_all_jobs( - ctx: &LibraryContext, - ) -> Result<(), prisma_client_rust::QueryError> { + pub async fn clear_all_jobs(ctx: &Library) -> Result<(), prisma_client_rust::QueryError> { ctx.db.job().delete_many(vec![]).exec().await?; invalidate_query!(ctx, "jobs.getHistory"); @@ -192,7 +190,7 @@ impl JobManager { } } - pub async fn resume_jobs(self: Arc, ctx: &LibraryContext) -> Result<(), JobError> { + pub async fn resume_jobs(self: Arc, ctx: &Library) -> Result<(), JobError> { let paused_jobs = ctx .db .job() @@ -261,7 +259,7 @@ impl JobManager { Ok(()) } - async fn dispatch_job(self: Arc, ctx: &LibraryContext, mut job: Box) { + async fn dispatch_job(self: Arc, ctx: &Library, mut job: Box) { // create worker to process job let mut running_workers = self.running_workers.write().await; if running_workers.len() < MAX_WORKERS { @@ -377,7 +375,7 @@ impl JobReport { } } - pub async fn create(&self, ctx: &LibraryContext) -> Result<(), JobError> { + pub async fn create(&self, ctx: &Library) -> Result<(), JobError> { ctx.db .job() .create( @@ -391,7 +389,7 @@ impl JobReport { .await?; Ok(()) } - pub async fn update(&self, ctx: &LibraryContext) -> Result<(), JobError> { + pub async fn update(&self, ctx: &Library) -> Result<(), JobError> { ctx.db .job() .update( diff --git a/core/src/job/worker.rs b/core/src/job/worker.rs index fbf908c6a..1e4849111 100644 --- a/core/src/job/worker.rs +++ b/core/src/job/worker.rs @@ -1,6 +1,6 @@ use crate::invalidate_query; use crate::job::{DynJob, JobError, JobManager, JobReportUpdate, JobStatus}; -use crate::library::LibraryContext; +use crate::library::Library; use std::{sync::Arc, time::Duration}; use tokio::sync::oneshot; use tokio::{ @@ -29,7 +29,7 @@ pub enum WorkerEvent { #[derive(Clone)] pub struct WorkerContext { - pub library_ctx: LibraryContext, + pub library: Library, events_tx: UnboundedSender, shutdown_tx: Arc>, } @@ -85,7 +85,7 @@ impl Worker { pub async fn spawn( job_manager: Arc, worker_mutex: Arc>, - ctx: LibraryContext, + library: Library, ) -> Result<(), JobError> { let mut worker = worker_mutex.lock().await; // we capture the worker receiver channel so state can be updated from inside the worker @@ -107,26 +107,25 @@ impl Worker { worker.report.status = JobStatus::Running; if matches!(old_status, JobStatus::Queued) { - worker.report.create(&ctx).await?; + worker.report.create(&library).await?; } else { - worker.report.update(&ctx).await?; + worker.report.update(&library).await?; } drop(worker); - invalidate_query!(ctx, "jobs.isRunning"); + invalidate_query!(library, "jobs.isRunning"); - let library_ctx = ctx.clone(); // spawn task to handle receiving events from the worker tokio::spawn(Worker::track_progress( Arc::clone(&worker_mutex), worker_events_rx, - library_ctx.clone(), + library.clone(), )); // spawn task to handle running the job tokio::spawn(async move { let worker_ctx = WorkerContext { - library_ctx, + library: library.clone(), events_tx: worker_events_tx, shutdown_tx: job_manager.shutdown_tx(), }; @@ -180,7 +179,7 @@ impl Worker { if let Err(e) = done_rx.await { error!("failed to wait for worker completion: {:#?}", e); } - job_manager.complete(&ctx, job_id, job_hash).await; + job_manager.complete(&library, job_id, job_hash).await; }); Ok(()) @@ -189,7 +188,7 @@ impl Worker { async fn track_progress( worker: Arc>, mut worker_events_rx: UnboundedReceiver, - library: LibraryContext, + library: Library, ) { let mut last = Instant::now(); diff --git a/core/src/lib.rs b/core/src/lib.rs index 3328a5e3d..8bfaabdaf 100644 --- a/core/src/lib.rs +++ b/core/src/lib.rs @@ -130,8 +130,8 @@ impl Node { .await?; // Adding already existing locations for location management - for library_ctx in library_manager.get_all_libraries_ctx().await { - for location in library_ctx + for library in library_manager.get_all_libraries().await { + for location in library .db .location() .find_many(vec![]) @@ -144,7 +144,7 @@ impl Node { ); vec![] }) { - if let Err(e) = location_manager.add(location.id, library_ctx.clone()).await { + if let Err(e) = location_manager.add(location.id, library.clone()).await { error!("Failed to add location to location manager: {:#?}", e); } } @@ -156,8 +156,8 @@ impl Node { let inner_library_manager = Arc::clone(&library_manager); let inner_jobs = Arc::clone(&jobs); tokio::spawn(async move { - for library_ctx in inner_library_manager.get_all_libraries_ctx().await { - if let Err(e) = Arc::clone(&inner_jobs).resume_jobs(&library_ctx).await { + for library in inner_library_manager.get_all_libraries().await { + if let Err(e) = Arc::clone(&inner_jobs).resume_jobs(&library).await { error!("Failed to resume jobs for library. {:#?}", e); } } diff --git a/core/src/library/library_config.rs b/core/src/library/config.rs similarity index 100% rename from core/src/library/library_config.rs rename to core/src/library/config.rs diff --git a/core/src/library/library_ctx.rs b/core/src/library/library.rs similarity index 96% rename from core/src/library/library_ctx.rs rename to core/src/library/library.rs index 7343a4816..9dbbe7a17 100644 --- a/core/src/library/library_ctx.rs +++ b/core/src/library/library.rs @@ -17,7 +17,7 @@ use super::LibraryConfig; /// LibraryContext holds context for a library which can be passed around the application. #[derive(Clone)] -pub struct LibraryContext { +pub struct Library { /// id holds the ID of the current library. pub id: Uuid, /// local_id holds the local ID of the current library. @@ -35,7 +35,7 @@ pub struct LibraryContext { pub(super) node_context: NodeContext, } -impl Debug for LibraryContext { +impl Debug for Library { fn fmt(&self, f: &mut Formatter<'_>) -> std::fmt::Result { // Rolling out this implementation because `NodeContext` contains a DynJob which is // troublesome to implement Debug trait @@ -48,7 +48,7 @@ impl Debug for LibraryContext { } } -impl LibraryContext { +impl Library { pub(crate) async fn spawn_job(&self, job: Box) { self.node_context.jobs.clone().ingest(self, job).await; } diff --git a/core/src/library/library_manager.rs b/core/src/library/manager.rs similarity index 97% rename from core/src/library/library_manager.rs rename to core/src/library/manager.rs index c1aabfad2..6c2facfcf 100644 --- a/core/src/library/library_manager.rs +++ b/core/src/library/manager.rs @@ -25,14 +25,14 @@ use tokio::sync::RwLock; use tracing::debug; use uuid::Uuid; -use super::{LibraryConfig, LibraryConfigWrapped, LibraryContext}; +use super::{Library, LibraryConfig, LibraryConfigWrapped}; /// LibraryManager is a singleton that manages all libraries for a node. pub struct LibraryManager { /// libraries_dir holds the path to the directory where libraries are stored. libraries_dir: PathBuf, /// libraries holds the list of libraries which are currently loaded into the node. - libraries: RwLock>, + libraries: RwLock>, /// node_context holds the context for the node which this library manager is running on. pub node_context: NodeContext, } @@ -227,7 +227,7 @@ impl LibraryManager { .collect() } - pub(crate) async fn get_all_libraries_ctx(&self) -> Vec { + pub(crate) async fn get_all_libraries(&self) -> Vec { self.libraries.read().await.clone() } @@ -282,7 +282,7 @@ impl LibraryManager { } // get_ctx will return the library context for the given library id. - pub(crate) async fn get_ctx(&self, library_id: Uuid) -> Option { + pub(crate) async fn get_ctx(&self, library_id: Uuid) -> Option { self.libraries .read() .await @@ -297,7 +297,7 @@ impl LibraryManager { db_path: impl AsRef, config: LibraryConfig, node_context: NodeContext, - ) -> Result { + ) -> Result { let db_path = db_path.as_ref(); let db = Arc::new( load_and_migrate(&format!( @@ -339,7 +339,7 @@ impl LibraryManager { let (sync_manager, _) = SyncManager::new(&db, id); - Ok(LibraryContext { + Ok(Library { id, local_id: node_data.id, config, diff --git a/core/src/library/mod.rs b/core/src/library/mod.rs index a9bb53fda..a93237e0f 100644 --- a/core/src/library/mod.rs +++ b/core/src/library/mod.rs @@ -1,7 +1,7 @@ -mod library_config; -mod library_ctx; -mod library_manager; +mod config; +mod library; +mod manager; -pub use library_config::*; -pub use library_ctx::*; -pub use library_manager::*; +pub use config::*; +pub use library::*; +pub use manager::*; diff --git a/core/src/location/file_path_helper.rs b/core/src/location/file_path_helper.rs index e27d63e08..30fc671fe 100644 --- a/core/src/location/file_path_helper.rs +++ b/core/src/location/file_path_helper.rs @@ -1,4 +1,4 @@ -use crate::{library::LibraryContext, prisma::file_path}; +use crate::{library::Library, prisma::file_path}; use std::sync::atomic::{AtomicI32, Ordering}; @@ -8,10 +8,10 @@ static LAST_FILE_PATH_ID: AtomicI32 = AtomicI32::new(0); file_path::select!(file_path_id_only { id }); -pub async fn get_max_file_path_id(library_ctx: &LibraryContext) -> Result { +pub async fn get_max_file_path_id(library: &Library) -> Result { let mut last_id = LAST_FILE_PATH_ID.load(Ordering::Acquire); if last_id == 0 { - last_id = fetch_max_file_path_id(library_ctx).await?; + last_id = fetch_max_file_path_id(library).await?; LAST_FILE_PATH_ID.store(last_id, Ordering::Release); } @@ -22,8 +22,8 @@ pub fn set_max_file_path_id(id: i32) { LAST_FILE_PATH_ID.store(id, Ordering::Relaxed); } -async fn fetch_max_file_path_id(library_ctx: &LibraryContext) -> Result { - Ok(library_ctx +async fn fetch_max_file_path_id(library: &Library) -> Result { + Ok(library .db .file_path() .find_first(vec![]) @@ -37,7 +37,7 @@ async fn fetch_max_file_path_id(library_ctx: &LibraryContext) -> Result) -> Result<(), JobError> { // grab the next id so we can increment in memory for batch inserting - let first_file_id = get_max_file_path_id(&ctx.library_ctx).await?; + let first_file_id = get_max_file_path_id(&ctx.library).await?; let mut indexer_rules_by_kind: HashMap> = HashMap::with_capacity(state.init.location.indexer_rules.len()); @@ -218,7 +218,7 @@ impl StatefulJob for IndexerJob { ctx: WorkerContext, state: &mut JobState, ) -> Result<(), JobError> { - let LibraryContext { sync, db, .. } = &ctx.library_ctx; + let Library { sync, db, .. } = &ctx.library; let location = &state.init.location; diff --git a/core/src/location/indexer/rules.rs b/core/src/location/indexer/rules.rs index 2ac10e923..189832a2c 100644 --- a/core/src/location/indexer/rules.rs +++ b/core/src/location/indexer/rules.rs @@ -1,5 +1,5 @@ use crate::{ - library::LibraryContext, + library::Library, location::indexer::IndexerError, prisma::{indexer_rule, PrismaClient}, }; @@ -29,7 +29,7 @@ pub struct IndexerRuleCreateArgs { } impl IndexerRuleCreateArgs { - pub async fn create(self, ctx: &LibraryContext) -> Result { + pub async fn create(self, ctx: &Library) -> Result { let parameters = match self.kind { RuleKind::AcceptFilesByGlob | RuleKind::RejectFilesByGlob => rmp_serde::to_vec( &Glob::new(&serde_json::from_slice::(&self.parameters)?)?, diff --git a/core/src/location/manager/helpers.rs b/core/src/location/manager/helpers.rs index 1343301da..8ed1961b4 100644 --- a/core/src/location/manager/helpers.rs +++ b/core/src/location/manager/helpers.rs @@ -1,4 +1,4 @@ -use crate::{library::LibraryContext, prisma::location}; +use crate::{library::Library, prisma::location}; use std::{ collections::{HashMap, HashSet}, @@ -17,17 +17,17 @@ type LocationAndLibraryKey = (LocationId, LibraryId); const LOCATION_CHECK_INTERVAL: Duration = Duration::from_secs(5); -pub(super) async fn check_online(location: &location::Data, library_ctx: &LibraryContext) -> bool { +pub(super) async fn check_online(location: &location::Data, library: &Library) -> bool { let pub_id = &location.pub_id; - if location.node_id == library_ctx.node_local_id { + if location.node_id == library.node_local_id { match fs::metadata(&location.path).await { Ok(_) => { - library_ctx.location_manager().add_online(pub_id).await; + library.location_manager().add_online(pub_id).await; true } Err(e) if e.kind() == ErrorKind::NotFound => { - library_ctx.location_manager().remove_online(pub_id).await; + library.location_manager().remove_online(pub_id).await; false } Err(e) => { @@ -37,17 +37,17 @@ pub(super) async fn check_online(location: &location::Data, library_ctx: &Librar } } else { // In this case, we don't have a `local_path`, but this location was marked as online - library_ctx.location_manager().remove_online(pub_id).await; + library.location_manager().remove_online(pub_id).await; false } } pub(super) async fn location_check_sleep( location_id: LocationId, - library_ctx: LibraryContext, -) -> (LocationId, LibraryContext) { + library: Library, +) -> (LocationId, Library) { sleep(LOCATION_CHECK_INTERVAL).await; - (location_id, library_ctx) + (location_id, library) } pub(super) fn watch_location( @@ -101,11 +101,8 @@ pub(super) fn drop_location( } } -pub(super) async fn get_location( - location_id: i32, - library_ctx: &LibraryContext, -) -> Option { - library_ctx +pub(super) async fn get_location(location_id: i32, library: &Library) -> Option { + library .db .location() .find_unique(location::id::equals(location_id)) @@ -138,28 +135,23 @@ pub(super) fn subtract_location_path( pub(super) async fn handle_remove_location_request( location_id: LocationId, - library_ctx: LibraryContext, + library: Library, response_tx: oneshot::Sender>, forced_unwatch: &mut HashSet, locations_watched: &mut HashMap, locations_unwatched: &mut HashMap, to_remove: &mut HashSet, ) { - let key = (location_id, library_ctx.id); - if let Some(location) = get_location(location_id, &library_ctx).await { - if location.node_id == library_ctx.node_local_id { - unwatch_location( - location, - library_ctx.id, - locations_watched, - locations_unwatched, - ); + let key = (location_id, library.id); + if let Some(location) = get_location(location_id, &library).await { + if location.node_id == library.node_local_id { + unwatch_location(location, library.id, locations_watched, locations_unwatched); locations_unwatched.remove(&key); forced_unwatch.remove(&key); } else { drop_location( location_id, - library_ctx.id, + library.id, "Dropping location from location manager, because we don't have a `local_path` anymore", locations_watched, locations_unwatched @@ -168,7 +160,7 @@ pub(super) async fn handle_remove_location_request( } else { drop_location( location_id, - library_ctx.id, + library.id, "Removing location from manager, as we failed to fetch from db", locations_watched, locations_unwatched, @@ -183,7 +175,7 @@ pub(super) async fn handle_remove_location_request( pub(super) async fn handle_stop_watcher_request( location_id: LocationId, - library_ctx: LibraryContext, + library: Library, response_tx: oneshot::Sender>, forced_unwatch: &mut HashSet, locations_watched: &mut HashMap, @@ -191,25 +183,20 @@ pub(super) async fn handle_stop_watcher_request( ) { async fn inner( location_id: LocationId, - library_ctx: LibraryContext, + library: Library, forced_unwatch: &mut HashSet, locations_watched: &mut HashMap, locations_unwatched: &mut HashMap, ) -> Result<(), LocationManagerError> { - let key = (location_id, library_ctx.id); + let key = (location_id, library.id); if !forced_unwatch.contains(&key) && locations_watched.contains_key(&key) { - get_location(location_id, &library_ctx) + get_location(location_id, &library) .await .ok_or_else(|| LocationManagerError::FailedToStopOrReinitWatcher { reason: String::from("failed to fetch location from db"), }) .map(|location| { - unwatch_location( - location, - library_ctx.id, - locations_watched, - locations_unwatched, - ); + unwatch_location(location, library.id, locations_watched, locations_unwatched); forced_unwatch.insert(key); }) } else { @@ -220,7 +207,7 @@ pub(super) async fn handle_stop_watcher_request( let _ = response_tx.send( inner( location_id, - library_ctx, + library, forced_unwatch, locations_watched, locations_unwatched, @@ -231,7 +218,7 @@ pub(super) async fn handle_stop_watcher_request( pub(super) async fn handle_reinit_watcher_request( location_id: LocationId, - library_ctx: LibraryContext, + library: Library, response_tx: oneshot::Sender>, forced_unwatch: &mut HashSet, locations_watched: &mut HashMap, @@ -239,25 +226,20 @@ pub(super) async fn handle_reinit_watcher_request( ) { async fn inner( location_id: LocationId, - library_ctx: LibraryContext, + library: Library, forced_unwatch: &mut HashSet, locations_watched: &mut HashMap, locations_unwatched: &mut HashMap, ) -> Result<(), LocationManagerError> { - let key = (location_id, library_ctx.id); + let key = (location_id, library.id); if forced_unwatch.contains(&key) && locations_unwatched.contains_key(&key) { - get_location(location_id, &library_ctx) + get_location(location_id, &library) .await .ok_or_else(|| LocationManagerError::FailedToStopOrReinitWatcher { reason: String::from("failed to fetch location from db"), }) .map(|location| { - watch_location( - location, - library_ctx.id, - locations_watched, - locations_unwatched, - ); + watch_location(location, library.id, locations_watched, locations_unwatched); forced_unwatch.remove(&key); }) } else { @@ -268,7 +250,7 @@ pub(super) async fn handle_reinit_watcher_request( let _ = response_tx.send( inner( location_id, - library_ctx, + library, forced_unwatch, locations_watched, locations_unwatched, @@ -279,14 +261,14 @@ pub(super) async fn handle_reinit_watcher_request( pub(super) fn handle_ignore_path_request( location_id: LocationId, - library_ctx: LibraryContext, + library: Library, path: PathBuf, ignore: bool, response_tx: oneshot::Sender>, locations_watched: &HashMap, ) { let _ = response_tx.send( - if let Some(watcher) = locations_watched.get(&(location_id, library_ctx.id)) { + if let Some(watcher) = locations_watched.get(&(location_id, library.id)) { watcher.ignore_path(path, ignore) } else { Ok(()) diff --git a/core/src/location/manager/mod.rs b/core/src/location/manager/mod.rs index 1cec8f982..c103d81d0 100644 --- a/core/src/location/manager/mod.rs +++ b/core/src/location/manager/mod.rs @@ -1,4 +1,4 @@ -use crate::library::LibraryContext; +use crate::library::Library; use std::{ collections::BTreeSet, @@ -39,7 +39,7 @@ enum ManagementMessageAction { #[allow(dead_code)] pub struct LocationManagementMessage { location_id: LocationId, - library_ctx: LibraryContext, + library: Library, action: ManagementMessageAction, response_tx: oneshot::Sender>, } @@ -56,7 +56,7 @@ enum WatcherManagementMessageAction { #[allow(dead_code)] pub struct WatcherManagementMessage { location_id: LocationId, - library_ctx: LibraryContext, + library: Library, action: WatcherManagementMessageAction, response_tx: oneshot::Sender>, } @@ -154,7 +154,7 @@ impl LocationManager { async fn location_management_message( &self, location_id: LocationId, - library_ctx: LibraryContext, + library: Library, action: ManagementMessageAction, ) -> Result<(), LocationManagerError> { #[cfg(feature = "location-watcher")] @@ -164,7 +164,7 @@ impl LocationManager { self.location_management_tx .send(LocationManagementMessage { location_id, - library_ctx, + library: library, action, response_tx: tx, }) @@ -182,7 +182,7 @@ impl LocationManager { async fn watcher_management_message( &self, location_id: LocationId, - library_ctx: LibraryContext, + library: Library, action: WatcherManagementMessageAction, ) -> Result<(), LocationManagerError> { #[cfg(feature = "location-watcher")] @@ -192,7 +192,7 @@ impl LocationManager { self.watcher_management_tx .send(WatcherManagementMessage { location_id, - library_ctx, + library: library, action, response_tx: tx, }) @@ -208,42 +208,38 @@ impl LocationManager { pub async fn add( &self, location_id: LocationId, - library_ctx: LibraryContext, + library: Library, ) -> Result<(), LocationManagerError> { - self.location_management_message(location_id, library_ctx, ManagementMessageAction::Add) + self.location_management_message(location_id, library, ManagementMessageAction::Add) .await } pub async fn remove( &self, location_id: LocationId, - library_ctx: LibraryContext, + library: Library, ) -> Result<(), LocationManagerError> { - self.location_management_message(location_id, library_ctx, ManagementMessageAction::Remove) + self.location_management_message(location_id, library, ManagementMessageAction::Remove) .await } pub async fn stop_watcher( &self, location_id: LocationId, - library_ctx: LibraryContext, + library: Library, ) -> Result<(), LocationManagerError> { - self.watcher_management_message( - location_id, - library_ctx, - WatcherManagementMessageAction::Stop, - ) - .await + self.watcher_management_message(location_id, library, WatcherManagementMessageAction::Stop) + .await } pub async fn reinit_watcher( &self, location_id: LocationId, - library_ctx: LibraryContext, + library: Library, ) -> Result<(), LocationManagerError> { self.watcher_management_message( location_id, - library_ctx, + library, WatcherManagementMessageAction::Reinit, ) .await @@ -252,13 +248,13 @@ impl LocationManager { pub async fn temporary_stop( &self, location_id: LocationId, - library_ctx: LibraryContext, + library: Library, ) -> Result { - self.stop_watcher(location_id, library_ctx.clone()).await?; + self.stop_watcher(location_id, library.clone()).await?; Ok(StopWatcherGuard { location_id, - library_ctx: Some(library_ctx), + library: Some(library), manager: self, }) } @@ -266,14 +262,14 @@ impl LocationManager { pub async fn temporary_ignore_events_for_path( &self, location_id: LocationId, - library_ctx: LibraryContext, + library: Library, path: impl AsRef, ) -> Result { let path = path.as_ref().to_path_buf(); self.watcher_management_message( location_id, - library_ctx.clone(), + library.clone(), WatcherManagementMessageAction::IgnoreEventsForPath { path: path.clone(), ignore: true, @@ -283,7 +279,7 @@ impl LocationManager { Ok(IgnoreEventsForPathGuard { location_id, - library_ctx: Some(library_ctx), + library: Some(library), manager: self, path: Some(path), }) @@ -319,7 +315,7 @@ impl LocationManager { // Location management messages Some(LocationManagementMessage{ location_id, - library_ctx, + library, action, response_tx }) = location_management_rx.recv() => { @@ -327,27 +323,27 @@ impl LocationManager { // To add a new location ManagementMessageAction::Add => { - if let Some(location) = get_location(location_id, &library_ctx).await { - let is_online = check_online(&location, &library_ctx).await; + if let Some(location) = get_location(location_id, &library).await { + let is_online = check_online(&location, &library).await; let _ = response_tx.send( - LocationWatcher::new(location, library_ctx.clone()) + LocationWatcher::new(location, library.clone()) .await .map(|mut watcher| { if is_online { watcher.watch(); locations_watched.insert( - (location_id, library_ctx.id), + (location_id, library.id), watcher ); } else { locations_unwatched.insert( - (location_id, library_ctx.id), + (location_id, library.id), watcher ); } to_check_futures.push( - location_check_sleep(location_id, library_ctx) + location_check_sleep(location_id, library) ); } ) @@ -364,7 +360,7 @@ impl LocationManager { ManagementMessageAction::Remove => { handle_remove_location_request( location_id, - library_ctx, + library, response_tx, &mut forced_unwatch, &mut locations_watched, @@ -378,7 +374,7 @@ impl LocationManager { // Watcher management messages Some(WatcherManagementMessage{ location_id, - library_ctx, + library, action, response_tx, }) = watcher_management_rx.recv() => { @@ -387,7 +383,7 @@ impl LocationManager { WatcherManagementMessageAction::Stop => { handle_stop_watcher_request( location_id, - library_ctx, + library, response_tx, &mut forced_unwatch, &mut locations_watched, @@ -399,7 +395,7 @@ impl LocationManager { WatcherManagementMessageAction::Reinit => { handle_reinit_watcher_request( location_id, - library_ctx, + library, response_tx, &mut forced_unwatch, &mut locations_watched, @@ -411,7 +407,7 @@ impl LocationManager { WatcherManagementMessageAction::IgnoreEventsForPath { path, ignore } => { handle_ignore_path_request( location_id, - library_ctx, + library, path, ignore, response_tx, @@ -422,36 +418,36 @@ impl LocationManager { } // Periodically checking locations - Some((location_id, library_ctx)) = to_check_futures.next() => { - let key = (location_id, library_ctx.id); + Some((location_id, library)) = to_check_futures.next() => { + let key = (location_id, library.id); if to_remove.contains(&key) { // The time to check came for an already removed library, so we just ignore it to_remove.remove(&key); - } else if let Some(location) = get_location(location_id, &library_ctx).await { - if location.node_id == library_ctx.node_local_id { - if check_online(&location, &library_ctx).await + } else if let Some(location) = get_location(location_id, &library).await { + if location.node_id == library.node_local_id { + if check_online(&location, &library).await && !forced_unwatch.contains(&key) { watch_location( location, - library_ctx.id, + library.id, &mut locations_watched, &mut locations_unwatched, ); } else { unwatch_location( location, - library_ctx.id, + library.id, &mut locations_watched, &mut locations_unwatched, ); } - to_check_futures.push(location_check_sleep(location_id, library_ctx)); + to_check_futures.push(location_check_sleep(location_id, library)); } else { drop_location( location_id, - library_ctx.id, + library.id, "Dropping location from location manager, because \ we don't have a `local_path` anymore", &mut locations_watched, @@ -462,7 +458,7 @@ impl LocationManager { } else { drop_location( location_id, - library_ctx.id, + library.id, "Removing location from manager, as we failed to fetch from db", &mut locations_watched, &mut locations_unwatched, @@ -524,7 +520,7 @@ impl Drop for LocationManager { pub struct StopWatcherGuard<'m> { manager: &'m LocationManager, location_id: LocationId, - library_ctx: Option, + library: Option, } impl Drop for StopWatcherGuard<'_> { @@ -533,7 +529,7 @@ impl Drop for StopWatcherGuard<'_> { // FIXME: change this Drop to async drop in the future if let Err(e) = block_on( self.manager - .reinit_watcher(self.location_id, self.library_ctx.take().unwrap()), + .reinit_watcher(self.location_id, self.library.take().unwrap()), ) { error!("Failed to reinit watcher on stop watcher guard drop: {e}"); } @@ -546,7 +542,7 @@ pub struct IgnoreEventsForPathGuard<'m> { manager: &'m LocationManager, path: Option, location_id: LocationId, - library_ctx: Option, + library: Option, } impl Drop for IgnoreEventsForPathGuard<'_> { @@ -555,7 +551,7 @@ impl Drop for IgnoreEventsForPathGuard<'_> { // FIXME: change this Drop to async drop in the future if let Err(e) = block_on(self.manager.watcher_management_message( self.location_id, - self.library_ctx.take().unwrap(), + self.library.take().unwrap(), WatcherManagementMessageAction::IgnoreEventsForPath { path: self.path.take().unwrap(), ignore: false, diff --git a/core/src/location/manager/watcher/linux.rs b/core/src/location/manager/watcher/linux.rs index 7eb5f8f54..1e347d648 100644 --- a/core/src/location/manager/watcher/linux.rs +++ b/core/src/location/manager/watcher/linux.rs @@ -1,5 +1,5 @@ use crate::{ - library::LibraryContext, + library::Library, location::{indexer::indexer_job::indexer_job_location, manager::LocationManagerError}, }; @@ -27,7 +27,7 @@ impl EventHandler for LinuxEventHandler { async fn handle_event( &mut self, location: indexer_job_location::Data, - library_ctx: &LibraryContext, + library: &Library, event: Event, ) -> Result<(), LocationManagerError> { trace!("Received Linux event: {:#?}", event); @@ -35,16 +35,16 @@ impl EventHandler for LinuxEventHandler { match event.kind { EventKind::Access(AccessKind::Close(AccessMode::Write)) => { // If a file was closed with write mode, then it was updated or created - file_creation_or_update(&location, &event, library_ctx).await?; + file_creation_or_update(&location, &event, library).await?; } EventKind::Create(CreateKind::Folder) => { - create_dir(&location, &event, library_ctx).await?; + create_dir(&location, &event, library).await?; } EventKind::Modify(ModifyKind::Name(RenameMode::Both)) => { - rename_both_event(&location, &event, library_ctx).await?; + rename_both_event(&location, &event, library).await?; } EventKind::Remove(remove_kind) => { - remove_event(&location, &event, remove_kind, library_ctx).await?; + remove_event(&location, &event, remove_kind, library).await?; } other_event_kind => { trace!("Other Linux event that we don't handle for now: {other_event_kind:#?}"); diff --git a/core/src/location/manager/watcher/macos.rs b/core/src/location/manager/watcher/macos.rs index 62f1d5039..1aa7625c6 100644 --- a/core/src/location/manager/watcher/macos.rs +++ b/core/src/location/manager/watcher/macos.rs @@ -1,5 +1,5 @@ use crate::{ - library::LibraryContext, + library::Library, location::{indexer::indexer_job::indexer_job_location, manager::LocationManagerError}, }; @@ -33,7 +33,7 @@ impl EventHandler for MacOsEventHandler { async fn handle_event( &mut self, location: indexer_job_location::Data, - library_ctx: &LibraryContext, + library: &Library, event: Event, ) -> Result<(), LocationManagerError> { trace!("Received MacOS event: {:#?}", event); @@ -50,12 +50,12 @@ impl EventHandler for MacOsEventHandler { } } - create_dir(&location, &event, library_ctx).await?; + create_dir(&location, &event, library).await?; self.latest_created_dir = Some(event); } EventKind::Modify(ModifyKind::Data(DataChange::Content)) => { // If a file had its content modified, then it was updated or created - file_creation_or_update(&location, &event, library_ctx).await?; + file_creation_or_update(&location, &event, library).await?; } EventKind::Modify(ModifyKind::Name(RenameMode::Any)) => { match self.rename_stack.take() { @@ -63,19 +63,13 @@ impl EventHandler for MacOsEventHandler { self.rename_stack = Some(event); } Some(from_event) => { - rename( - &event.paths[0], - &from_event.paths[0], - &location, - library_ctx, - ) - .await?; + rename(&event.paths[0], &from_event.paths[0], &location, library).await?; } } } EventKind::Remove(remove_kind) => { - remove_event(&location, &event, remove_kind, library_ctx).await?; + remove_event(&location, &event, remove_kind, library).await?; } other_event_kind => { trace!("Other MacOS event that we don't handle for now: {other_event_kind:#?}"); diff --git a/core/src/location/manager/watcher/mod.rs b/core/src/location/manager/watcher/mod.rs index a70f470cc..91a3c2e42 100644 --- a/core/src/location/manager/watcher/mod.rs +++ b/core/src/location/manager/watcher/mod.rs @@ -1,5 +1,5 @@ use crate::{ - library::LibraryContext, + library::Library, prisma::{file_path, location}, }; @@ -53,7 +53,7 @@ trait EventHandler { async fn handle_event( &mut self, location: indexer_job_location::Data, - library_ctx: &LibraryContext, + library: &Library, event: Event, ) -> Result<(), LocationManagerError>; } @@ -70,7 +70,7 @@ pub(super) struct LocationWatcher { impl LocationWatcher { pub(super) async fn new( location: location::Data, - library_ctx: LibraryContext, + library: Library, ) -> Result { let (events_tx, events_rx) = mpsc::unbounded_channel(); let (ignore_path_tx, ignore_path_rx) = mpsc::unbounded_channel(); @@ -97,7 +97,7 @@ impl LocationWatcher { let handle = tokio::spawn(Self::handle_watch_events( location.id, - library_ctx, + library, events_rx, ignore_path_rx, stop_rx, @@ -114,7 +114,7 @@ impl LocationWatcher { async fn handle_watch_events( location_id: LocationId, - library_ctx: LibraryContext, + library: Library, mut events_rx: mpsc::UnboundedReceiver>, mut ignore_path_rx: mpsc::UnboundedReceiver, mut stop_rx: oneshot::Receiver<()>, @@ -132,7 +132,7 @@ impl LocationWatcher { location_id, event, &mut event_handler, - &library_ctx, + &library, &paths_to_ignore, ).await { error!("Failed to handle location file system event: \ @@ -166,14 +166,14 @@ impl LocationWatcher { location_id: LocationId, event: Event, event_handler: &mut impl EventHandler, - library_ctx: &LibraryContext, + library: &Library, ignore_paths: &HashSet, ) -> Result<(), LocationManagerError> { if !check_event(&event, ignore_paths) { return Ok(()); } - let Some(location) = fetch_location(library_ctx, location_id) + let Some(location) = fetch_location(library, location_id) .include(indexer_job_location::include()) .exec() .await? @@ -182,18 +182,12 @@ impl LocationWatcher { return Ok(()); }; - if !library_ctx - .location_manager() - .is_online(&location.pub_id) - .await - { + if !library.location_manager().is_online(&location.pub_id).await { warn!("Tried to handle event for offline location: "); return Ok(()); } - event_handler - .handle_event(location, library_ctx, event) - .await + event_handler.handle_event(location, library, event).await } pub(super) fn ignore_path( diff --git a/core/src/location/manager/watcher/utils.rs b/core/src/location/manager/watcher/utils.rs index ce2bc1f58..48ff69a66 100644 --- a/core/src/location/manager/watcher/utils.rs +++ b/core/src/location/manager/watcher/utils.rs @@ -1,6 +1,6 @@ use crate::{ invalidate_query, - library::LibraryContext, + library::Library, location::{ delete_directory, file_path_helper::create_file_path, @@ -49,9 +49,9 @@ pub(super) fn check_event(event: &Event, ignore_paths: &HashSet) -> boo pub(super) async fn create_dir( location: &indexer_job_location::Data, event: &Event, - library_ctx: &LibraryContext, + library: &Library, ) -> Result<(), LocationManagerError> { - if location.node_id != library_ctx.node_local_id { + if location.node_id != library.node_local_id { return Ok(()); } @@ -65,7 +65,7 @@ pub(super) async fn create_dir( return Ok(()); }; - let parent_directory = get_parent_dir(location.id, &subpath, library_ctx).await?; + let parent_directory = get_parent_dir(location.id, &subpath, library).await?; trace!("parent_directory: {:?}", parent_directory); @@ -75,7 +75,7 @@ pub(super) async fn create_dir( }; let created_path = create_file_path( - library_ctx, + library, location.id, subpath .to_str() @@ -94,7 +94,7 @@ pub(super) async fn create_dir( info!("Created path: {}", created_path.materialized_path); - invalidate_query!(library_ctx, "locations.getExplorerData"); + invalidate_query!(library, "locations.getExplorerData"); Ok(()) } @@ -102,9 +102,9 @@ pub(super) async fn create_dir( pub(super) async fn create_file( location: &indexer_job_location::Data, event: &Event, - library_ctx: &LibraryContext, + library: &Library, ) -> Result<(), LocationManagerError> { - if location.node_id != library_ctx.node_local_id { + if location.node_id != library.node_local_id { return Ok(()); } @@ -114,19 +114,19 @@ pub(super) async fn create_file( event.paths[0].display() ); - let db = &library_ctx.db; + let db = &library.db; let Some(materialized_path) = subtract_location_path(&location.path, &event.paths[0]) else { return Ok(()) }; let Some(parent_directory) = - get_parent_dir(location.id, &materialized_path, library_ctx).await? + get_parent_dir(location.id, &materialized_path, library).await? else { warn!("Watcher found a path without parent"); return Ok(()) }; let created_file = create_file_path( - library_ctx, + library, location.id, materialized_path .to_str() @@ -209,16 +209,10 @@ pub(super) async fn create_file( trace!("object: {:#?}", object); if !object.has_thumbnail && !created_file.extension.is_empty() { - generate_thumbnail( - &created_file.extension, - &cas_id, - &event.paths[0], - library_ctx, - ) - .await; + generate_thumbnail(&created_file.extension, &cas_id, &event.paths[0], library).await; } - invalidate_query!(library_ctx, "locations.getExplorerData"); + invalidate_query!(library, "locations.getExplorerData"); Ok(()) } @@ -226,29 +220,29 @@ pub(super) async fn create_file( pub(super) async fn file_creation_or_update( location: &indexer_job_location::Data, event: &Event, - library_ctx: &LibraryContext, + library: &Library, ) -> Result<(), LocationManagerError> { if let Some(ref file_path) = - get_existing_file_path(location, &event.paths[0], false, library_ctx).await? + get_existing_file_path(location, &event.paths[0], false, library).await? { - inner_update_file(location, file_path, event, library_ctx).await + inner_update_file(location, file_path, event, library).await } else { // We received None because it is a new file - create_file(location, event, library_ctx).await + create_file(location, event, library).await } } pub(super) async fn update_file( location: &indexer_job_location::Data, event: &Event, - library_ctx: &LibraryContext, + library: &Library, ) -> Result<(), LocationManagerError> { - if location.node_id == library_ctx.node_local_id { + if location.node_id == library.node_local_id { if let Some(ref file_path) = - get_existing_file_path(location, &event.paths[0], false, library_ctx).await? + get_existing_file_path(location, &event.paths[0], false, library).await? { - let ret = inner_update_file(location, file_path, event, library_ctx).await; - invalidate_query!(library_ctx, "locations.getExplorerData"); + let ret = inner_update_file(location, file_path, event, library).await; + invalidate_query!(library, "locations.getExplorerData"); ret } else { Err(LocationManagerError::UpdateNonExistingFile( @@ -264,7 +258,7 @@ async fn inner_update_file( location: &indexer_job_location::Data, file_path: &file_path_with_object::Data, event: &Event, - library_ctx: &LibraryContext, + library: &Library, ) -> Result<(), LocationManagerError> { trace!( "Location: updating file: {}", @@ -281,7 +275,7 @@ async fn inner_update_file( if let Some(old_cas_id) = &file_path.cas_id { if old_cas_id != &cas_id { // file content changed - library_ctx + library .db .file_path() .update( @@ -314,14 +308,14 @@ async fn inner_update_file( { // if this file had a thumbnail previously, we update it to match the new content if !file_path.extension.is_empty() { - generate_thumbnail(&file_path.extension, &cas_id, &event.paths[0], library_ctx) + generate_thumbnail(&file_path.extension, &cas_id, &event.paths[0], library) .await; } } } } - invalidate_query!(library_ctx, "locations.getExplorerData"); + invalidate_query!(library, "locations.getExplorerData"); Ok(()) } @@ -329,16 +323,16 @@ async fn inner_update_file( pub(super) async fn rename_both_event( location: &indexer_job_location::Data, event: &Event, - library_ctx: &LibraryContext, + library: &Library, ) -> Result<(), LocationManagerError> { - rename(&event.paths[1], &event.paths[0], location, library_ctx).await + rename(&event.paths[1], &event.paths[0], location, library).await } pub(super) async fn rename( new_path: impl AsRef, old_path: impl AsRef, location: &indexer_job_location::Data, - library_ctx: &LibraryContext, + library: &Library, ) -> Result<(), LocationManagerError> { let mut old_path_materialized = extract_materialized_path(location, old_path.as_ref())? .to_str() @@ -351,8 +345,7 @@ pub(super) async fn rename( .expect("Found non-UTF-8 path") .to_string(); - if let Some(file_path) = get_existing_file_or_directory(location, old_path, library_ctx).await? - { + if let Some(file_path) = get_existing_file_or_directory(location, old_path, library).await? { // If the renamed path is a directory, we have to update every successor if file_path.is_dir { if !old_path_materialized.ends_with('/') { @@ -362,7 +355,7 @@ pub(super) async fn rename( new_path_materialized_str += "/"; } - let updated = library_ctx + let updated = library .db ._execute_raw( raw!( @@ -377,7 +370,7 @@ pub(super) async fn rename( trace!("Updated {updated} file_paths"); } - library_ctx + library .db .file_path() .update( @@ -406,7 +399,7 @@ pub(super) async fn rename( ) .exec() .await?; - invalidate_query!(library_ctx, "locations.getExplorerData"); + invalidate_query!(library, "locations.getExplorerData"); } Ok(()) @@ -416,13 +409,13 @@ pub(super) async fn remove_event( location: &indexer_job_location::Data, event: &Event, remove_kind: RemoveKind, - library_ctx: &LibraryContext, + library: &Library, ) -> Result<(), LocationManagerError> { trace!("removed {remove_kind:#?}"); // if it doesn't either way, then we don't care if let Some(file_path) = - get_existing_file_or_directory(location, &event.paths[0], library_ctx).await? + get_existing_file_or_directory(location, &event.paths[0], library).await? { // check file still exists on disk match fs::metadata(&event.paths[0]).await { @@ -432,10 +425,10 @@ pub(super) async fn remove_event( Err(e) if e.kind() == ErrorKind::NotFound => { // if is doesn't, we can remove it safely from our db if file_path.is_dir { - delete_directory(library_ctx, location.id, Some(file_path.materialized_path)) + delete_directory(library, location.id, Some(file_path.materialized_path)) .await?; } else { - library_ctx + library .db .file_path() .delete(file_path::location_id_id(location.id, file_path.id)) @@ -443,7 +436,7 @@ pub(super) async fn remove_event( .await?; if let Some(object_id) = file_path.object_id { - library_ctx + library .db .object() .delete_many(vec![ @@ -459,7 +452,7 @@ pub(super) async fn remove_event( Err(e) => return Err(e.into()), } - invalidate_query!(library_ctx, "locations.getExplorerData"); + invalidate_query!(library, "locations.getExplorerData"); } Ok(()) @@ -481,7 +474,7 @@ async fn get_existing_file_path( location: &indexer_job_location::Data, path: impl AsRef, is_dir: bool, - library_ctx: &LibraryContext, + library: &Library, ) -> Result, LocationManagerError> { let mut materialized_path = extract_materialized_path(location, path)? .to_str() @@ -491,7 +484,7 @@ async fn get_existing_file_path( materialized_path += "/"; } - library_ctx + library .db .file_path() .find_first(vec![file_path::materialized_path::equals( @@ -507,14 +500,13 @@ async fn get_existing_file_path( async fn get_existing_file_or_directory( location: &indexer_job_location::Data, path: impl AsRef, - library_ctx: &LibraryContext, + library: &Library, ) -> Result, LocationManagerError> { let mut maybe_file_path = - get_existing_file_path(location, path.as_ref(), false, library_ctx).await?; + get_existing_file_path(location, path.as_ref(), false, library).await?; // First we just check if this path was a file in our db, if it isn't then we check for a directory if maybe_file_path.is_none() { - maybe_file_path = - get_existing_file_path(location, path.as_ref(), true, library_ctx).await?; + maybe_file_path = get_existing_file_path(location, path.as_ref(), true, library).await?; } Ok(maybe_file_path) @@ -523,7 +515,7 @@ async fn get_existing_file_or_directory( async fn get_parent_dir( location_id: LocationId, path: impl AsRef, - library_ctx: &LibraryContext, + library: &Library, ) -> Result, LocationManagerError> { let mut parent_path_str = path .as_ref() @@ -539,7 +531,7 @@ async fn get_parent_dir( parent_path_str += "/"; } - library_ctx + library .db .file_path() .find_first(vec![ @@ -555,10 +547,10 @@ async fn generate_thumbnail( extension: &str, cas_id: &str, file_path: impl AsRef, - library_ctx: &LibraryContext, + library: &Library, ) { let file_path = file_path.as_ref(); - let output_path = library_ctx + let output_path = library .config() .data_directory() .join(THUMBNAIL_CACHE_DIR_NAME) diff --git a/core/src/location/manager/watcher/windows.rs b/core/src/location/manager/watcher/windows.rs index 719eb9d30..879636ef0 100644 --- a/core/src/location/manager/watcher/windows.rs +++ b/core/src/location/manager/watcher/windows.rs @@ -1,5 +1,5 @@ use crate::{ - library::LibraryContext, + library::Library, location::{indexer::indexer_job::indexer_job_location, manager::LocationManagerError}, }; @@ -34,7 +34,7 @@ impl EventHandler for WindowsEventHandler { async fn handle_event( &mut self, location: indexer_job_location::Data, - library_ctx: &LibraryContext, + library: &Library, event: Event, ) -> Result<(), LocationManagerError> { trace!("Received Windows event: {:#?}", event); @@ -45,16 +45,16 @@ impl EventHandler for WindowsEventHandler { if metadata.is_file() { self.create_file_stack = Some(event); } else { - create_dir(&location, &event, library_ctx).await?; + create_dir(&location, &event, library).await?; } } EventKind::Modify(ModifyKind::Any) => { let metadata = fs::metadata(&event.paths[0]).await?; if metadata.is_file() { if let Some(create_file_event) = self.create_file_stack.take() { - create_file(&location, &create_file_event, library_ctx).await?; + create_file(&location, &create_file_event, library).await?; } else { - update_file(&location, &event, library_ctx).await?; + update_file(&location, &event, library).await?; } } else { warn!("Unexpected Windows modify event on a directory"); @@ -68,16 +68,10 @@ impl EventHandler for WindowsEventHandler { .rename_stack .take() .expect("Unexpectedly missing rename from windows event"); - rename( - &event.paths[0], - &from_event.paths[0], - &location, - library_ctx, - ) - .await?; + rename(&event.paths[0], &from_event.paths[0], &location, library).await?; } EventKind::Remove(remove_kind) => { - remove_event(&location, &event, remove_kind, library_ctx).await?; + remove_event(&location, &event, remove_kind, library).await?; } other_event_kind => { diff --git a/core/src/location/mod.rs b/core/src/location/mod.rs index 79c075ddf..b2bd42749 100644 --- a/core/src/location/mod.rs +++ b/core/src/location/mod.rs @@ -1,7 +1,7 @@ use crate::{ invalidate_query, job::Job, - library::LibraryContext, + library::Library, object::{ identifier_job::full_identifier_job::{FullFileIdentifierJob, FullFileIdentifierJobInit}, preview::{ThumbnailJob, ThumbnailJobInit}, @@ -45,10 +45,7 @@ pub struct LocationCreateArgs { } impl LocationCreateArgs { - pub async fn create( - self, - ctx: &LibraryContext, - ) -> Result { + pub async fn create(self, ctx: &Library) -> Result { let path_metadata = match fs::metadata(&self.path).await { Ok(metadata) => metadata, Err(e) if e.kind() == io::ErrorKind::NotFound => { @@ -107,7 +104,7 @@ impl LocationCreateArgs { pub async fn add_library( self, - ctx: &LibraryContext, + ctx: &Library, ) -> Result { let mut metadata = SpacedriveLocationMetadataFile::try_load(&self.path) .await? @@ -163,8 +160,8 @@ pub struct LocationUpdateArgs { } impl LocationUpdateArgs { - pub async fn update(self, ctx: &LibraryContext) -> Result<(), LocationError> { - let LibraryContext { sync, db, .. } = &ctx; + pub async fn update(self, ctx: &Library) -> Result<(), LocationError> { + let Library { sync, db, .. } = &ctx; let location = fetch_location(ctx, self.id) .include(location::include!({ indexer_rules })) @@ -265,14 +262,14 @@ impl LocationUpdateArgs { } } -pub fn fetch_location(ctx: &LibraryContext, location_id: i32) -> location::FindUnique { +pub fn fetch_location(ctx: &Library, location_id: i32) -> location::FindUnique { ctx.db .location() .find_unique(location::id::equals(location_id)) } async fn link_location_and_indexer_rules( - ctx: &LibraryContext, + ctx: &Library, location_id: i32, rules_ids: &[i32], ) -> Result<(), LocationError> { @@ -291,7 +288,7 @@ async fn link_location_and_indexer_rules( } pub async fn scan_location( - ctx: &LibraryContext, + ctx: &Library, location: indexer_job_location::Data, ) -> Result<(), LocationError> { if location.node_id != ctx.node_local_id { @@ -324,10 +321,10 @@ pub async fn scan_location( } pub async fn relink_location( - ctx: &LibraryContext, + ctx: &Library, location_path: impl AsRef, ) -> Result<(), LocationError> { - let LibraryContext { db, id, sync, .. } = &ctx; + let Library { db, id, sync, .. } = &ctx; let mut metadata = SpacedriveLocationMetadataFile::try_load(&location_path) .await? @@ -362,12 +359,12 @@ pub async fn relink_location( } async fn create_location( - ctx: &LibraryContext, + ctx: &Library, location_pub_id: Uuid, location_path: impl AsRef, indexer_rules_ids: &[i32], ) -> Result { - let LibraryContext { db, sync, .. } = &ctx; + let Library { db, sync, .. } = &ctx; let location_path = location_path.as_ref(); @@ -425,8 +422,8 @@ async fn create_location( Ok(location) } -pub async fn delete_location(ctx: &LibraryContext, location_id: i32) -> Result<(), LocationError> { - let LibraryContext { db, .. } = ctx; +pub async fn delete_location(ctx: &Library, location_id: i32) -> Result<(), LocationError> { + let Library { db, .. } = ctx; ctx.location_manager() .remove(location_id, ctx.clone()) @@ -466,7 +463,7 @@ file_path::select!(file_path_object_id_only { object_id }); /// Will delete a directory recursively with Objects if left as orphans /// this function is used to delete a location and when ingesting directory deletion events pub async fn delete_directory( - ctx: &LibraryContext, + ctx: &Library, location_id: i32, parent_materialized_path: Option, ) -> Result<(), QueryError> { @@ -517,13 +514,13 @@ pub async fn delete_directory( // check if a path exists in our database at that location // pub async fn check_virtual_path_exists( -// library_ctx: &LibraryContext, +// library: &Library, // location_id: i32, // subpath: impl AsRef, // ) -> Result { // let path = subpath.as_ref().to_str().unwrap().to_string(); -// let file_path = library_ctx +// let file_path = library // .db // .file_path() // .find_first(vec![ diff --git a/core/src/object/fs/copy.rs b/core/src/object/fs/copy.rs index 6b68a1542..f0b9380c2 100644 --- a/core/src/object/fs/copy.rs +++ b/core/src/object/fs/copy.rs @@ -62,14 +62,14 @@ impl StatefulJob for FileCopierJob { async fn init(&self, ctx: WorkerContext, state: &mut JobState) -> Result<(), JobError> { let source_fs_info = context_menu_fs_info( - &ctx.library_ctx.db, + &ctx.library.db, state.init.source_location_id, state.init.source_path_id, ) .await?; let mut full_target_path = - get_path_from_location_id(&ctx.library_ctx.db, state.init.target_location_id).await?; + get_path_from_location_id(&ctx.library.db, state.init.target_location_id).await?; // add the currently viewed subdirectory to the location root full_target_path.push(&state.init.target_path); diff --git a/core/src/object/fs/create.rs b/core/src/object/fs/create.rs index 9b2d0a555..81fe58364 100644 --- a/core/src/object/fs/create.rs +++ b/core/src/object/fs/create.rs @@ -13,16 +13,16 @@ // location_id: i32, // path: &str, // name: Option<&str>, -// library_ctx: &LibraryContext, +// library: &LibraryContext, // ) -> Result<(), VirtualFSError> { -// let location = fetch_location(library_ctx, location_id) +// let location = fetch_location(library, location_id) // .exec() // .await? // .ok_or(LocationError::IdNotFound(location_id))?; // let name = name.unwrap_or("Untitled Folder"); -// let exists = check_virtual_path_exists(library_ctx, location_id, subpath).await?; +// let exists = check_virtual_path_exists(library, location_id, subpath).await?; // std::fs::create_dir_all(&obj_path)?; diff --git a/core/src/object/fs/cut.rs b/core/src/object/fs/cut.rs index f51b1426f..20e58307e 100644 --- a/core/src/object/fs/cut.rs +++ b/core/src/object/fs/cut.rs @@ -41,14 +41,14 @@ impl StatefulJob for FileCutterJob { async fn init(&self, ctx: WorkerContext, state: &mut JobState) -> Result<(), JobError> { let source_fs_info = context_menu_fs_info( - &ctx.library_ctx.db, + &ctx.library.db, state.init.source_location_id, state.init.source_path_id, ) .await?; let mut full_target_path = - get_path_from_location_id(&ctx.library_ctx.db, state.init.target_location_id).await?; + get_path_from_location_id(&ctx.library.db, state.init.target_location_id).await?; full_target_path.push(&state.init.target_path); state.steps = [FileCutterJobStep { diff --git a/core/src/object/fs/decrypt.rs b/core/src/object/fs/decrypt.rs index 3dbcc2a80..11c2f26ce 100644 --- a/core/src/object/fs/decrypt.rs +++ b/core/src/object/fs/decrypt.rs @@ -45,12 +45,9 @@ impl StatefulJob for FileDecryptorJob { async fn init(&self, ctx: WorkerContext, state: &mut JobState) -> Result<(), JobError> { // enumerate files to decrypt // populate the steps with them (local file paths) - let fs_info = context_menu_fs_info( - &ctx.library_ctx.db, - state.init.location_id, - state.init.path_id, - ) - .await?; + let fs_info = + context_menu_fs_info(&ctx.library.db, state.init.location_id, state.init.path_id) + .await?; state.steps = VecDeque::new(); state.steps.push_back(FileDecryptorJobStep { fs_info }); @@ -67,7 +64,7 @@ impl StatefulJob for FileDecryptorJob { ) -> Result<(), JobError> { let step = &state.steps[0]; let info = &step.fs_info; - let key_manager = &ctx.library_ctx.key_manager; + let key_manager = &ctx.library.key_manager; // handle overwriting checks, and making sure there's enough available space let output_path = state.init.output_path.clone().map_or_else( diff --git a/core/src/object/fs/delete.rs b/core/src/object/fs/delete.rs index df2460363..dc1daf0b5 100644 --- a/core/src/object/fs/delete.rs +++ b/core/src/object/fs/delete.rs @@ -31,12 +31,9 @@ impl StatefulJob for FileDeleterJob { } async fn init(&self, ctx: WorkerContext, state: &mut JobState) -> Result<(), JobError> { - let fs_info = context_menu_fs_info( - &ctx.library_ctx.db, - state.init.location_id, - state.init.path_id, - ) - .await?; + let fs_info = + context_menu_fs_info(&ctx.library.db, state.init.location_id, state.init.path_id) + .await?; state.steps = [fs_info].into_iter().collect(); diff --git a/core/src/object/fs/encrypt.rs b/core/src/object/fs/encrypt.rs index b5b096fdb..06ccfc958 100644 --- a/core/src/object/fs/encrypt.rs +++ b/core/src/object/fs/encrypt.rs @@ -1,4 +1,4 @@ -use crate::{job::*, library::LibraryContext}; +use crate::{job::*, library::Library}; use std::path::PathBuf; @@ -58,15 +58,12 @@ impl StatefulJob for FileEncryptorJob { } async fn init(&self, ctx: WorkerContext, state: &mut JobState) -> Result<(), JobError> { - let step = context_menu_fs_info( - &ctx.library_ctx.db, - state.init.location_id, - state.init.path_id, - ) - .await - .map_err(|_| JobError::MissingData { - value: String::from("file_path that matches both location id and path id"), - })?; + let step = + context_menu_fs_info(&ctx.library.db, state.init.location_id, state.init.path_id) + .await + .map_err(|_| JobError::MissingData { + value: String::from("file_path that matches both location id and path id"), + })?; state.steps = [step].into_iter().collect(); @@ -82,7 +79,7 @@ impl StatefulJob for FileEncryptorJob { ) -> Result<(), JobError> { let info = &state.steps[0]; - let LibraryContext { key_manager, .. } = &ctx.library_ctx; + let Library { key_manager, .. } = &ctx.library; if !info.path_data.is_dir { // handle overwriting checks, and making sure there's enough available space @@ -120,11 +117,11 @@ impl StatefulJob for FileEncryptorJob { )?; let _guard = ctx - .library_ctx + .library .location_manager() .temporary_ignore_events_for_path( state.init.location_id, - ctx.library_ctx.clone(), + ctx.library.clone(), &output_path, ) .await?; @@ -181,7 +178,7 @@ impl StatefulJob for FileEncryptorJob { // may not be the best - pvm isn't guaranteed to be webp let pvm_path = ctx - .library_ctx + .library .config() .data_directory() .join("thumbnails") diff --git a/core/src/object/fs/erase.rs b/core/src/object/fs/erase.rs index 2b1c66d28..a5e799d67 100644 --- a/core/src/object/fs/erase.rs +++ b/core/src/object/fs/erase.rs @@ -55,12 +55,9 @@ impl StatefulJob for FileEraserJob { } async fn init(&self, ctx: WorkerContext, state: &mut JobState) -> Result<(), JobError> { - let fs_info = context_menu_fs_info( - &ctx.library_ctx.db, - state.init.location_id, - state.init.path_id, - ) - .await?; + let fs_info = + context_menu_fs_info(&ctx.library.db, state.init.location_id, state.init.path_id) + .await?; state.data = Some(fs_info.clone()); diff --git a/core/src/object/identifier_job/full_identifier_job.rs b/core/src/object/identifier_job/full_identifier_job.rs index e69f82830..43a087388 100644 --- a/core/src/object/identifier_job/full_identifier_job.rs +++ b/core/src/object/identifier_job/full_identifier_job.rs @@ -1,7 +1,7 @@ use crate::{ invalidate_query, job::{JobError, JobReportUpdate, JobResult, JobState, StatefulJob, WorkerContext}, - library::LibraryContext, + library::Library, prisma::{file_path, location}, }; @@ -68,7 +68,7 @@ impl StatefulJob for FullFileIdentifierJob { info!("Identifying orphan File Paths..."); let location_id = state.init.location_id; - let db = &ctx.library_ctx.db; + let db = &ctx.library.db; let location = db .location() @@ -77,7 +77,7 @@ impl StatefulJob for FullFileIdentifierJob { .await? .ok_or(IdentifierJobError::MissingLocation(state.init.location_id))?; - let orphan_count = count_orphan_file_paths(&ctx.library_ctx, location_id).await?; + let orphan_count = count_orphan_file_paths(&ctx.library, location_id).await?; info!("Found {} orphan file paths", orphan_count); let task_count = (orphan_count as f64 / CHUNK_SIZE as f64).ceil() as usize; @@ -127,7 +127,7 @@ impl StatefulJob for FullFileIdentifierJob { // get chunk of orphans to process let file_paths = - get_orphan_file_paths(&ctx.library_ctx, &data.cursor, data.location.id).await?; + get_orphan_file_paths(&ctx.library, &data.cursor, data.location.id).await?; // if no file paths found, abort entire job early, there is nothing to do // if we hit this error, there is something wrong with the data/query @@ -147,7 +147,7 @@ impl StatefulJob for FullFileIdentifierJob { ); let (total_objects_created, total_objects_linked) = - identifier_job_step(&ctx.library_ctx, &data.location, &file_paths).await?; + identifier_job_step(&ctx.library, &data.location, &file_paths).await?; data.report.total_objects_created += total_objects_created; data.report.total_objects_linked += total_objects_linked; @@ -166,7 +166,7 @@ impl StatefulJob for FullFileIdentifierJob { )), ]); - invalidate_query!(ctx.library_ctx, "locations.getExplorerData"); + invalidate_query!(ctx.library, "locations.getExplorerData"); // let _remaining = count_orphan_file_paths(&ctx.core_ctx, location_id.into()).await?; Ok(()) @@ -198,7 +198,7 @@ fn orphan_path_filters(location_id: i32, file_path_id: Option) -> Vec Result { Ok(ctx @@ -214,7 +214,7 @@ async fn count_orphan_file_paths( } async fn get_orphan_file_paths( - ctx: &LibraryContext, + ctx: &Library, cursor: &FilePathIdAndLocationIdCursor, location_id: i32, ) -> Result, prisma_client_rust::QueryError> { diff --git a/core/src/object/identifier_job/mod.rs b/core/src/object/identifier_job/mod.rs index 270e76dea..e798f3d21 100644 --- a/core/src/object/identifier_job/mod.rs +++ b/core/src/object/identifier_job/mod.rs @@ -1,6 +1,6 @@ use crate::{ job::JobError, - library::LibraryContext, + library::Library, object::cas::generate_cas_id, prisma::{file_path, location, object, PrismaClient}, sync, @@ -76,7 +76,7 @@ impl FileMetadata { } async fn identifier_job_step( - LibraryContext { db, sync, .. }: &LibraryContext, + Library { db, sync, .. }: &Library, location: &location::Data, file_paths: &[file_path::Data], ) -> Result<(usize, usize), JobError> { diff --git a/core/src/object/preview/thumb.rs b/core/src/object/preview/thumb.rs index f1a1cfe06..6087a18c0 100644 --- a/core/src/object/preview/thumb.rs +++ b/core/src/object/preview/thumb.rs @@ -2,7 +2,7 @@ use crate::{ api::CoreEvent, invalidate_query, job::{JobError, JobReportUpdate, JobResult, JobState, StatefulJob, WorkerContext}, - library::LibraryContext, + library::Library, prisma::{file_path, location}, }; @@ -76,10 +76,10 @@ impl StatefulJob for ThumbnailJob { } async fn init(&self, ctx: WorkerContext, state: &mut JobState) -> Result<(), JobError> { - let LibraryContext { db, .. } = &ctx.library_ctx; + let Library { db, .. } = &ctx.library; let thumbnail_dir = ctx - .library_ctx + .library .config() .data_directory() .join(THUMBNAIL_CACHE_DIR_NAME); @@ -126,7 +126,7 @@ impl StatefulJob for ThumbnailJob { // query database for all image files in this location that need thumbnails let image_files = get_files_by_extensions( - &ctx.library_ctx, + &ctx.library, state.init.location_id, parent_directory_id, &sd_file_ext::extensions::ALL_IMAGE_EXTENSIONS @@ -144,7 +144,7 @@ impl StatefulJob for ThumbnailJob { let all_files = { // query database for all video files in this location that need thumbnails let video_files = get_files_by_extensions( - &ctx.library_ctx, + &ctx.library, state.init.location_id, parent_directory_id, &sd_file_ext::extensions::ALL_VIDEO_EXTENSIONS @@ -253,7 +253,7 @@ impl StatefulJob for ThumbnailJob { // // media_data::pixel_height::set(Some(media_data.height)), // ]; // let _ = ctx - // .library_ctx() + // .library() // .db // .media_data() // .upsert( @@ -268,13 +268,13 @@ impl StatefulJob for ThumbnailJob { } if !state.init.background { - ctx.library_ctx.emit(CoreEvent::NewThumbnail { + ctx.library.emit(CoreEvent::NewThumbnail { cas_id: cas_id.clone(), }); }; // With this invalidate query, we update the user interface to show each new thumbnail - invalidate_query!(ctx.library_ctx, "locations.getExplorerData"); + invalidate_query!(ctx.library, "locations.getExplorerData"); } else { info!("Thumb exists, skipping... {}", output_path.display()); } @@ -346,7 +346,7 @@ pub async fn generate_video_thumbnail>( } async fn get_files_by_extensions( - ctx: &LibraryContext, + ctx: &Library, location_id: i32, _parent_file_path_id: i32, extensions: &[Extension], diff --git a/core/src/object/validation/validator_job.rs b/core/src/object/validation/validator_job.rs index 7d29d0d09..a5438ddf4 100644 --- a/core/src/object/validation/validator_job.rs +++ b/core/src/object/validation/validator_job.rs @@ -5,7 +5,7 @@ use std::{collections::VecDeque, path::PathBuf}; use crate::{ job::{JobError, JobReportUpdate, JobResult, JobState, StatefulJob, WorkerContext}, - library::LibraryContext, + library::Library, prisma::{file_path, location}, sync, }; @@ -60,7 +60,7 @@ impl StatefulJob for ObjectValidatorJob { } async fn init(&self, ctx: WorkerContext, state: &mut JobState) -> Result<(), JobError> { - let db = &ctx.library_ctx.db; + let db = &ctx.library.db; state.steps = db .file_path() @@ -97,7 +97,7 @@ impl StatefulJob for ObjectValidatorJob { ctx: WorkerContext, state: &mut JobState, ) -> Result<(), JobError> { - let LibraryContext { db, sync, .. } = &ctx.library_ctx; + let Library { db, sync, .. } = &ctx.library; let file_path = &state.steps[0]; let data = state.data.as_ref().expect("fatal: missing job state"); diff --git a/core/src/sync/manager.rs b/core/src/sync/manager.rs index 740be014a..88788c587 100644 --- a/core/src/sync/manager.rs +++ b/core/src/sync/manager.rs @@ -140,31 +140,16 @@ impl SyncManager { } pub async fn get_ops(&self) -> prisma_client_rust::Result> { - let db = &self.db; - - let owned = db - .owned_operation() - .find_many(vec![]) - .include(owned_operation::include!({ node })) - .exec() - .await? - .into_iter() - .flat_map(|op| { - Some(CRDTOperation { - id: Uuid::from_slice(&op.id).ok()?, - node: Uuid::from_slice(&op.node.pub_id).ok()?, - timestamp: NTP64(op.timestamp as u64), - typ: CRDTOperationType::Owned(OwnedOperation { - model: op.model, - items: serde_json::from_slice(&op.data).ok()?, - }), - }) - }); - - let shared = db + Ok(self + .db .shared_operation() .find_many(vec![]) - .include(shared_operation::include!({ node })) + .order_by(shared_operation::timestamp::order( + prisma_client_rust::Direction::Asc, + )) + .include(shared_operation::include!({ node: select { + pub_id + } })) .exec() .await? .into_iter() @@ -179,13 +164,8 @@ impl SyncManager { data: serde_json::from_slice(&op.data).ok()?, }), }) - }); - - let mut result: Vec = owned.chain(shared).collect(); - - result.sort_by(|a, b| a.timestamp.cmp(&b.timestamp)); - - Ok(result) + }) + .collect()) } pub async fn ingest_op(&self, op: CRDTOperation) -> prisma_client_rust::Result<()> { diff --git a/core/src/volume.rs b/core/src/volume.rs index ec58fa3a1..b6846420f 100644 --- a/core/src/volume.rs +++ b/core/src/volume.rs @@ -1,4 +1,4 @@ -use crate::{library::LibraryContext, prisma::volume::*}; +use crate::{library::Library, prisma::volume::*}; use rspc::Type; use serde::{Deserialize, Serialize}; @@ -38,7 +38,7 @@ impl From for rspc::Error { } } -pub async fn save_volume(ctx: &LibraryContext) -> Result<(), VolumeError> { +pub async fn save_volume(ctx: &Library) -> Result<(), VolumeError> { let volumes = get_volumes()?; // enter all volumes associate with this client add to db