use std::{ collections::BTreeMap, path::PathBuf, sync::{Arc, RwLock}, thread, time::Duration, }; use fuser::INodeNo; use sea_orm::{DatabaseConnection, EntityTrait}; use tokio::sync::Notify; use tokio_stream::StreamExt; use crate::item::Item; use crate::origins::network::transport::NetworkTransport; use crate::origins::{FileWatcher, WatcherHandle}; use crate::proto::ManifestEntry as ProtoManifestEntry; use tracing::{info, warn}; const INITIAL_RETRY_DELAY: Duration = Duration::from_secs(1); const MAX_RETRY_DELAY: Duration = Duration::from_secs(30); pub struct NetworkOriginFileWatcher { transport: NetworkTransport, runtime_handle: tokio::runtime::Handle, client: DatabaseConnection, destination: PathBuf, latest_manifest: Arc>>, } impl NetworkOriginFileWatcher { #[allow(clippy::too_many_arguments)] pub fn new( transport: NetworkTransport, runtime_handle: tokio::runtime::Handle, client: DatabaseConnection, destination: PathBuf, latest_manifest: Arc>>, ) -> Self { return NetworkOriginFileWatcher { transport, runtime_handle, client, destination, latest_manifest, }; } } impl FileWatcher for NetworkOriginFileWatcher { fn watch(&self, files: Arc>>) -> WatcherHandle { let runtime_handle = self.runtime_handle.clone(); let shutdown = Arc::new(Notify::new()); let state = WatcherState { transport: self.transport.clone(), client: self.client.clone(), destination: self.destination.clone(), latest_manifest: self.latest_manifest.clone(), files, }; let loop_shutdown = shutdown.clone(); let join = thread::spawn(move || { runtime_handle.block_on(state.run_loop(loop_shutdown)); }); return WatcherHandle::new(move || { shutdown.notify_one(); let _ = join.join(); }); } } #[derive(Clone)] struct WatcherState { transport: NetworkTransport, client: DatabaseConnection, destination: PathBuf, latest_manifest: Arc>>, files: Arc>>, } impl WatcherState { async fn run_loop(self, shutdown: Arc) { let mut retry_delay = INITIAL_RETRY_DELAY; loop { let subscribed = tokio::select! { biased; _ = shutdown.notified() => return, s = self.transport.subscribe_events() => s, }; match subscribed { Ok(mut stream) => { info!("network watcher: subscribed to /events"); retry_delay = INITIAL_RETRY_DELAY; loop { let item = tokio::select! { biased; _ = shutdown.notified() => return, item = stream.next() => item, }; match item { Some(Ok(_event)) => { if let Err(e) = self.reconcile_once().await { warn!(error = %e, "network watcher: reconcile after event failed"); } } Some(Err(e)) => { warn!(error = %e, "network watcher: stream error; reconnecting"); break; } None => break, } } } Err(e) => { warn!(error = %e, "network watcher: subscribe failed; will retry"); } } tokio::select! { biased; _ = shutdown.notified() => return, _ = tokio::time::sleep(retry_delay) => {} } retry_delay = (retry_delay * 2).min(MAX_RETRY_DELAY); if let Err(e) = self.reconcile_once().await { warn!(error = %e, "network watcher: poll reconcile failed"); } } } async fn reconcile_once(&self) -> anyhow::Result<()> { let client_entries: Vec<(u64, u64)> = crate::db::entities::Entity::find() .all(&self.client) .await? .into_iter() .map(|m| (m.inode as u64, m.hash as u64)) .collect(); let response = self.transport.reconcile(client_entries).await?; let mut changed_or_deleted: Vec = response.changed.iter().map(|ih| ih.inode as i64).collect(); changed_or_deleted.extend(response.deleted.iter().map(|i| *i as i64)); crate::db::cache::delete_cached_bytes_for(&changed_or_deleted, &self.client).await; let wanted: Vec = response.changed.iter().map(|ih| ih.inode).collect(); let mut current_manifest = self.latest_manifest.read().unwrap().clone(); if !wanted.is_empty() { let entries = self.transport.get_metadata(wanted).await?; for entry in entries { current_manifest.insert(entry.id, entry); } } for inode in &response.deleted { current_manifest.remove(inode); } *self.latest_manifest.write().unwrap() = current_manifest.clone(); let new_snapshot = super::build_snapshot_from_manifest(¤t_manifest, &self.destination, &self.client) .await?; { let mut files = self.files.lock().unwrap(); files.retain(|ino, _| new_snapshot.contains_key(ino)); for (ino, new_item) in &new_snapshot { match files.get(ino) { Some(existing) if existing.hash == new_item.hash => {} _ => { files.insert(*ino, new_item.clone()); } } } } info!( changed = response.changed.len(), deleted = response.deleted.len(), total = current_manifest.len(), "network watcher: reconcile applied" ); return Ok(()); } }