rename LibraryContext to Library (#595)

This commit is contained in:
Brendan Allan authored and GitHub committed 2023-03-07 01:15:20 -08:00
1 parent 31df51501e
commit 23ad538678
37 files changed
+375 -473

No files matched your search

+63 -75
View File
@@ -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<String>,
}
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(())
})
})
}
+2 -2
View File
@@ -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()
+2 -2
View File
@@ -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<Ctx> {
}
t(|_, mut args: LocationExplorerArgs, library| async move {
let LibraryContext { db, .. } = &library;
let Library { db, .. } = &library;
let location = db
.location()
+4 -4
View File
@@ -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()
+6 -6
View File
@@ -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),
))
)
+9 -9
View File
@@ -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<TBuiltResolver>,
) -> 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<TUnbuiltResultMarker> + Send,
TArg: DeserializeOwned + specta::Type + Send + 'static;
@@ -49,8 +49,8 @@ pub trait LibraryRequest {
) -> BuiltProcedureBuilder<TBuiltResolver>,
) -> 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<TUnbuiltResultMarker> + Send,
TArg: DeserializeOwned + specta::Type + Send + 'static;
@@ -79,8 +79,8 @@ where
) -> BuiltProcedureBuilder<TBuiltResolver>,
) -> 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<TUnbuiltResultMarker> + Send,
TArg: DeserializeOwned + specta::Type + Send + 'static,
{
@@ -125,8 +125,8 @@ where
) -> BuiltProcedureBuilder<TBuiltResolver>,
) -> 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<TUnbuiltResultMarker> + Send,
TArg: DeserializeOwned + specta::Type + Send + 'static,
{
+10 -12
View File
@@ -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<dyn DynJob>),
IngestJob(Library, Box<dyn DynJob>),
}
/// JobManager handles queueing and executing jobs using the `DynJob`
@@ -83,7 +83,7 @@ impl JobManager {
this
}
pub async fn ingest(self: Arc<Self>, ctx: &LibraryContext, job: Box<dyn DynJob>) {
pub async fn ingest(self: Arc<Self>, ctx: &Library, job: Box<dyn DynJob>) {
let job_hash = job.hash();
debug!(
"Ingesting job: <name='{}', hash='{}'>",
@@ -119,7 +119,7 @@ impl JobManager {
}
}
pub async fn complete(self: Arc<Self>, ctx: &LibraryContext, job_id: Uuid, job_hash: u64) {
pub async fn complete(self: Arc<Self>, 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<Vec<JobReport>, 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<Self>, ctx: &LibraryContext) -> Result<(), JobError> {
pub async fn resume_jobs(self: Arc<Self>, ctx: &Library) -> Result<(), JobError> {
let paused_jobs = ctx
.db
.job()
@@ -261,7 +259,7 @@ impl JobManager {
Ok(())
}
async fn dispatch_job(self: Arc<Self>, ctx: &LibraryContext, mut job: Box<dyn DynJob>) {
async fn dispatch_job(self: Arc<Self>, ctx: &Library, mut job: Box<dyn DynJob>) {
// 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(
+10 -11
View File
@@ -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<WorkerEvent>,
shutdown_tx: Arc<broadcast::Sender<()>>,
}
@@ -85,7 +85,7 @@ impl Worker {
pub async fn spawn(
job_manager: Arc<JobManager>,
worker_mutex: Arc<Mutex<Self>>,
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<Mutex<Self>>,
mut worker_events_rx: UnboundedReceiver<WorkerEvent>,
library: LibraryContext,
library: Library,
) {
let mut last = Instant::now();
+5 -5
View File
@@ -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);
}
}
File renamed without changes.
@@ -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<dyn DynJob>) {
self.node_context.jobs.clone().ingest(self, job).await;
}
@@ -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<Vec<LibraryContext>>,
libraries: RwLock<Vec<Library>>,
/// 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<LibraryContext> {
pub(crate) async fn get_all_libraries(&self) -> Vec<Library> {
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<LibraryContext> {
pub(crate) async fn get_ctx(&self, library_id: Uuid) -> Option<Library> {
self.libraries
.read()
.await
@@ -297,7 +297,7 @@ impl LibraryManager {
db_path: impl AsRef<Path>,
config: LibraryConfig,
node_context: NodeContext,
) -> Result<LibraryContext, LibraryManagerError> {
) -> Result<Library, LibraryManagerError> {
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,
+6 -6
View File
@@ -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::*;
+8 -8
View File
@@ -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<i32, QueryError> {
pub async fn get_max_file_path_id(library: &Library) -> Result<i32, QueryError> {
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<i32, QueryError> {
Ok(library_ctx
async fn fetch_max_file_path_id(library: &Library) -> Result<i32, QueryError> {
Ok(library
.db
.file_path()
.find_first(vec![])
@@ -37,7 +37,7 @@ async fn fetch_max_file_path_id(library_ctx: &LibraryContext) -> Result<i32, Que
#[cfg(feature = "location-watcher")]
pub async fn create_file_path(
library_ctx: &LibraryContext,
library: &Library,
location_id: i32,
mut materialized_path: String,
name: String,
@@ -49,7 +49,7 @@ pub async fn create_file_path(
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?;
}
// If this new file_path is a directory, materialized_path must end with "/"
@@ -59,7 +59,7 @@ pub async fn create_file_path(
let next_id = last_id + 1;
let created_path = library_ctx
let created_path = library
.db
.file_path()
.create(
+3 -3
View File
@@ -1,6 +1,6 @@
use crate::{
job::{JobError, JobReportUpdate, JobResult, JobState, StatefulJob, WorkerContext},
library::LibraryContext,
library::Library,
location::indexer::rules::RuleKind,
prisma::{file_path, location},
sync,
@@ -111,7 +111,7 @@ impl StatefulJob for IndexerJob {
/// Creates a vector of valid path buffers from a directory, chunked into batches of `BATCH_SIZE`.
async fn init(&self, ctx: WorkerContext, state: &mut JobState<Self>) -> 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<RuleKind, Vec<IndexerRule>> =
HashMap::with_capacity(state.init.location.indexer_rules.len());
@@ -218,7 +218,7 @@ impl StatefulJob for IndexerJob {
ctx: WorkerContext,
state: &mut JobState<Self>,
) -> Result<(), JobError> {
let LibraryContext { sync, db, .. } = &ctx.library_ctx;
let Library { sync, db, .. } = &ctx.library;
let location = &state.init.location;
+2 -2
View File
@@ -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<indexer_rule::Data, IndexerError> {
pub async fn create(self, ctx: &Library) -> Result<indexer_rule::Data, IndexerError> {
let parameters = match self.kind {
RuleKind::AcceptFilesByGlob | RuleKind::RejectFilesByGlob => rmp_serde::to_vec(
&Glob::new(&serde_json::from_slice::<String>(&self.parameters)?)?,
+32 -50
View File
@@ -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<location::Data> {
library_ctx
pub(super) async fn get_location(location_id: i32, library: &Library) -> Option<location::Data> {
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<Result<(), LocationManagerError>>,
forced_unwatch: &mut HashSet<LocationAndLibraryKey>,
locations_watched: &mut HashMap<LocationAndLibraryKey, LocationWatcher>,
locations_unwatched: &mut HashMap<LocationAndLibraryKey, LocationWatcher>,
to_remove: &mut HashSet<LocationAndLibraryKey>,
) {
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<Result<(), LocationManagerError>>,
forced_unwatch: &mut HashSet<LocationAndLibraryKey>,
locations_watched: &mut HashMap<LocationAndLibraryKey, LocationWatcher>,
@@ -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<LocationAndLibraryKey>,
locations_watched: &mut HashMap<LocationAndLibraryKey, LocationWatcher>,
locations_unwatched: &mut HashMap<LocationAndLibraryKey, LocationWatcher>,
) -> 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<Result<(), LocationManagerError>>,
forced_unwatch: &mut HashSet<LocationAndLibraryKey>,
locations_watched: &mut HashMap<LocationAndLibraryKey, LocationWatcher>,
@@ -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<LocationAndLibraryKey>,
locations_watched: &mut HashMap<LocationAndLibraryKey, LocationWatcher>,
locations_unwatched: &mut HashMap<LocationAndLibraryKey, LocationWatcher>,
) -> 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<Result<(), LocationManagerError>>,
locations_watched: &HashMap<LocationAndLibraryKey, LocationWatcher>,
) {
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(())
+48 -52
View File
@@ -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<Result<(), LocationManagerError>>,
}
@@ -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<Result<(), LocationManagerError>>,
}
@@ -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<StopWatcherGuard, LocationManagerError> {
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<Path>,
) -> Result<IgnoreEventsForPathGuard, LocationManagerError> {
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<LibraryContext>,
library: Option<Library>,
}
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<PathBuf>,
location_id: LocationId,
library_ctx: Option<LibraryContext>,
library: Option<Library>,
}
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,
+6 -6
View File
@@ -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:#?}");
+6 -12
View File
@@ -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:#?}");
+10 -16
View File
@@ -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<Self, LocationManagerError> {
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<notify::Result<Event>>,
mut ignore_path_rx: mpsc::UnboundedReceiver<IgnorePath>,
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<PathBuf>,
) -> 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: <id='{location_id}'>");
return Ok(());
}
event_handler
.handle_event(location, library_ctx, event)
.await
event_handler.handle_event(location, library, event).await
}
pub(super) fn ignore_path(
+48 -56
View File
@@ -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<PathBuf>) -> 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: <root_path ='{}'> 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<Path>,
old_path: impl AsRef<Path>,
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<Path>,
is_dir: bool,
library_ctx: &LibraryContext,
library: &Library,
) -> Result<Option<file_path_with_object::Data>, 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<Path>,
library_ctx: &LibraryContext,
library: &Library,
) -> Result<Option<file_path_with_object::Data>, 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<Path>,
library_ctx: &LibraryContext,
library: &Library,
) -> Result<Option<file_path::Data>, 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<Path>,
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)
+7 -13
View File
@@ -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 => {
+17 -20
View File
@@ -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<indexer_job_location::Data, LocationError> {
pub async fn create(self, ctx: &Library) -> Result<indexer_job_location::Data, LocationError> {
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<indexer_job_location::Data, LocationError> {
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<Path>,
) -> 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<Path>,
indexer_rules_ids: &[i32],
) -> Result<indexer_job_location::Data, LocationError> {
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<String>,
) -> 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<Path>,
// ) -> Result<bool, LocationError> {
// 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![
+2 -2
View File
@@ -62,14 +62,14 @@ impl StatefulJob for FileCopierJob {
async fn init(&self, ctx: WorkerContext, state: &mut JobState<Self>) -> 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);
+3 -3
View File
@@ -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)?;
+2 -2
View File
@@ -41,14 +41,14 @@ impl StatefulJob for FileCutterJob {
async fn init(&self, ctx: WorkerContext, state: &mut JobState<Self>) -> 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 {
+4 -7
View File
@@ -45,12 +45,9 @@ impl StatefulJob for FileDecryptorJob {
async fn init(&self, ctx: WorkerContext, state: &mut JobState<Self>) -> 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(
+3 -6
View File
@@ -31,12 +31,9 @@ impl StatefulJob for FileDeleterJob {
}
async fn init(&self, ctx: WorkerContext, state: &mut JobState<Self>) -> 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();
+11 -14
View File
@@ -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<Self>) -> 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")
+3 -6
View File
@@ -55,12 +55,9 @@ impl StatefulJob for FileEraserJob {
}
async fn init(&self, ctx: WorkerContext, state: &mut JobState<Self>) -> 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());
@@ -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<i32>) -> Vec<file_
}
async fn count_orphan_file_paths(
ctx: &LibraryContext,
ctx: &Library,
location_id: i32,
) -> Result<usize, prisma_client_rust::QueryError> {
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<Vec<file_path::Data>, prisma_client_rust::QueryError> {
+2 -2
View File
@@ -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> {
+9 -9
View File
@@ -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<Self>) -> 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<P: AsRef<Path>>(
}
async fn get_files_by_extensions(
ctx: &LibraryContext,
ctx: &Library,
location_id: i32,
_parent_file_path_id: i32,
extensions: &[Extension],
+3 -3
View File
@@ -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<Self>) -> 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<Self>,
) -> 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");
+10 -30
View File
@@ -140,31 +140,16 @@ impl SyncManager {
}
pub async fn get_ops(&self) -> prisma_client_rust::Result<Vec<CRDTOperation>> {
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<CRDTOperation> = 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<()> {
+2 -2
View File
@@ -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<VolumeError> 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