use crate::model::action::Action;
use crate::model::download_progress_item::DownloadProgressItem;
use crate::model::error::{LocalError, S3Error};
use crate::model::has_children::flatten_items;
use crate::model::local_data_item::LocalDataItem;
use crate::model::local_selected_item::LocalSelectedItem;
use crate::model::s3_data_item::S3DataItem;
use crate::model::s3_selected_item::S3SelectedItem;
use crate::model::state::{ActivePage, State};
use crate::model::transfer_state::TransferState;
use crate::model::upload_progress_item::UploadProgressItem;
use crate::services::local_data_fetcher::LocalDataFetcher;
use crate::services::s3_data_fetcher::S3DataFetcher;
use crate::services::task_registry::TaskRegistry;
use crate::services::transfer_manager::{JobId, TransferManager};
use crate::services::transfer_persistence::TransferPersistence;
use crate::services::transfer_state::TransferStateStore;
use crate::settings::file_credentials::FileCredential;
use crate::termination::{Interrupted, Terminator};
use crate::utils::get_data_dir;
use color_eyre::eyre;
use std::collections::HashMap;
use std::path::Path;
use std::sync::Arc;
use std::time::Duration;
use tokio::sync::mpsc::{Receiver, Sender, UnboundedReceiver, UnboundedSender};
use tokio::sync::{broadcast, mpsc, Mutex};
static S3_OPERATIONS_CONCURRENCY_LEVEL: usize = 8;
const STATE_CHANNEL_CAPACITY: usize = 100;
const PROGRESS_CHANNEL_CAPACITY: usize = 1000;
#[derive(Debug, Clone)]
pub enum TransferResult {
UploadComplete(LocalSelectedItem),
UploadFailed(LocalSelectedItem, String),
DownloadComplete(S3SelectedItem),
DownloadFailed(S3SelectedItem, String),
}
pub struct StateStore {
state_tx: Sender<State>,
task_registry: Arc<TaskRegistry>,
transfer_state_store: Option<Arc<TransferStateStore>>,
}
impl StateStore {
pub fn new() -> (Self, Receiver<State>) {
let (state_tx, state_rx) = mpsc::channel::<State>(STATE_CHANNEL_CAPACITY);
(
StateStore {
state_tx,
task_registry: Arc::new(TaskRegistry::new()),
transfer_state_store: None,
},
state_rx,
)
}
}
impl StateStore {
async fn enqueue_downloads(
transfer_manager: &TransferManager,
s3_selected_items: Vec<S3SelectedItem>,
job_item_map: &Arc<Mutex<HashMap<JobId, TransferJobInfo>>>,
) -> Vec<S3SelectedItem> {
let items_with_children = flatten_items(s3_selected_items);
let mut updated_items = Vec::new();
for mut item in items_with_children {
if item.job_id.is_some() || item.is_bucket || item.is_directory {
continue;
}
let s3_path = item.path.clone().unwrap_or_else(|| item.name.clone());
let local_path = format!("{}/{}", item.destination_dir, item.name);
let job_id = transfer_manager
.enqueue_download(s3_path.clone(), local_path, None)
.await;
item.job_id = Some(job_id);
item.transfer_state = TransferState::Pending;
job_item_map.lock().await.insert(
job_id,
TransferJobInfo::Download {
item: item.clone(),
},
);
updated_items.push(item);
}
updated_items
}
async fn enqueue_uploads(
transfer_manager: &TransferManager,
local_selected_items: Vec<LocalSelectedItem>,
job_item_map: &Arc<Mutex<HashMap<JobId, TransferJobInfo>>>,
) -> Vec<LocalSelectedItem> {
let items_with_children = flatten_items(local_selected_items);
let mut updated_items = Vec::new();
for mut item in items_with_children {
if item.job_id.is_some() || item.is_directory {
continue;
}
let local_path = item.path.clone();
let s3_path = if item.destination_path.is_empty() {
item.name.clone()
} else {
format!("{}/{}", item.destination_path, item.name)
};
let job_id = transfer_manager
.enqueue_upload(local_path.clone(), s3_path, None)
.await;
item.job_id = Some(job_id);
item.transfer_state = TransferState::Pending;
job_item_map.lock().await.insert(
job_id,
TransferJobInfo::Upload { item: item.clone() },
);
updated_items.push(item);
}
updated_items
}
async fn spawn_transfer_worker(
task_registry: &TaskRegistry,
transfer_manager: Arc<TransferManager>,
s3_data_fetcher: S3DataFetcher,
job_item_map: Arc<Mutex<HashMap<JobId, TransferJobInfo>>>,
result_tx: UnboundedSender<TransferResult>,
upload_progress_tx: Sender<UploadProgressItem>,
download_progress_tx: Sender<DownloadProgressItem>,
) {
task_registry.spawn_tracked("transfer-worker", async move {
loop {
if let Some((job, pause_signal)) = transfer_manager.try_get_next().await {
let job_id = job.id;
let job_info = job_item_map.lock().await.get(&job_id).cloned();
match job_info {
Some(TransferJobInfo::Upload { item }) => {
let up_tx = upload_progress_tx.clone();
let fetcher = s3_data_fetcher.clone();
let tm = transfer_manager.clone();
let tx = result_tx.clone();
tokio::spawn(async move {
match fetcher.upload_item(item.clone(), up_tx, Some(pause_signal)).await {
Ok(_) => {
tm.mark_completed(job_id).await;
let mut completed_item = item;
completed_item.transfer_state = TransferState::Completed;
let _ = tx.send(TransferResult::UploadComplete(completed_item));
}
Err(e) => {
let error_msg = e.to_string();
if !error_msg.contains("paused") {
tm.mark_failed(job_id, error_msg.clone()).await;
let mut failed_item = item;
failed_item.transfer_state =
TransferState::Failed(error_msg.clone());
let _ =
tx.send(TransferResult::UploadFailed(failed_item, error_msg));
}
}
}
});
}
Some(TransferJobInfo::Download { item }) => {
let down_tx = download_progress_tx.clone();
let fetcher = s3_data_fetcher.clone();
let tm = transfer_manager.clone();
let tx = result_tx.clone();
tokio::spawn(async move {
match fetcher.download_item(item.clone(), down_tx, Some(pause_signal)).await {
Ok(_) => {
tm.mark_completed(job_id).await;
let mut completed_item = item;
completed_item.transfer_state = TransferState::Completed;
let _ =
tx.send(TransferResult::DownloadComplete(completed_item));
}
Err(e) => {
let error_msg = e.to_string();
if !error_msg.contains("paused") {
tm.mark_failed(job_id, error_msg.clone()).await;
let mut failed_item = item;
failed_item.transfer_state =
TransferState::Failed(error_msg.clone());
let _ = tx
.send(TransferResult::DownloadFailed(failed_item, error_msg));
}
}
}
});
}
None => {
tracing::warn!("Job info not found for job_id: {}", job_id);
}
}
} else {
tokio::time::sleep(Duration::from_millis(100)).await;
}
}
}).await;
}
async fn fetch_s3_data(
&self,
bucket: Option<String>,
prefix: Option<String>,
s3_data_fetcher: S3DataFetcher,
s3_tx: UnboundedSender<(Option<String>, Option<String>, Vec<S3DataItem>)>,
) {
let task_name = match (&bucket, &prefix) {
(Some(b), Some(p)) => format!("fetch-s3:{}/{}", b, p),
(Some(b), None) => format!("fetch-s3:{}", b),
_ => "fetch-s3:buckets".to_string(),
};
self.task_registry.spawn_tracked(task_name, async move {
match s3_data_fetcher
.list_current_location(bucket.clone(), prefix.clone())
.await
{
Ok(data) => {
let _ = s3_tx.send((bucket.clone(), prefix.clone(), data));
}
Err(e) => {
tracing::error!("Failed to fetch S3 data: {}", e);
}
}
}).await;
}
async fn list_s3_data_recursive(
&self,
item: S3SelectedItem,
s3_data_fetcher: S3DataFetcher,
s3_full_list_tx: UnboundedSender<(Option<String>, Option<String>, Vec<S3DataItem>)>,
) {
let task_name = format!("list-s3-recursive:{}", item.name);
self.task_registry.spawn_tracked(task_name, async move {
let bucket_name = if item.is_bucket {
item.name
} else {
item.bucket.unwrap_or(item.name)
};
let path = if item.is_bucket {
None
} else {
item.path.clone()
};
match s3_data_fetcher
.list_all_objects(&bucket_name, path.clone())
.await
{
Ok(data) => {
let _ = s3_full_list_tx.send((Some(bucket_name), path.clone(), data));
}
Err(e) => {
tracing::error!("Failed to fetch S3 data: {}", e);
}
}
}).await;
}
async fn fetch_local_data(
&self,
dir_path: Option<String>,
local_data_fetcher: LocalDataFetcher,
local_tx: UnboundedSender<(String, Vec<LocalDataItem>)>,
) {
let path = Self::get_directory_path(dir_path);
let task_name = format!("fetch-local:{}", path.as_deref().unwrap_or("/"));
self.task_registry.spawn_tracked(task_name, async move {
match local_data_fetcher.read_directory(path.clone()).await {
Ok(data) => {
let _ = local_tx.send((path.clone().unwrap_or("/".to_string()), data));
}
Err(e) => {
tracing::error!("Failed to fetch local data: {}", e);
}
}
}).await;
}
fn get_directory_path(input_path: Option<String>) -> Option<String> {
match input_path {
Some(path) => {
let path = Path::new(&path);
if path.is_dir() {
path.to_str().map(String::from)
} else {
path.parent().and_then(|p| p.to_str().map(String::from))
}
}
None => None,
}
}
async fn move_back_local_data(
&self,
current_path: String,
local_data_fetcher: LocalDataFetcher,
local_tx: UnboundedSender<(String, Vec<LocalDataItem>)>,
) {
self.task_registry.spawn_tracked("move-back-local", async move {
let path = Path::new(¤t_path);
match local_data_fetcher.read_parent_directory().await {
Ok(data) => {
let _ = match path.parent() {
Some(p_path) => local_tx.send((p_path.to_string_lossy().to_string(), data)),
None => local_tx.send((current_path, data)),
};
}
Err(e) => {
tracing::error!("Failed to fetch local data: {}", e);
}
}
}).await;
}
async fn delete_local_data(
&self,
item: LocalSelectedItem,
local_data_fetcher: LocalDataFetcher,
local_deleted_tx: UnboundedSender<Option<LocalError>>,
) {
let path = item.path.clone();
let task_name = format!("delete-local:{}", item.name);
if item.is_directory {
self.task_registry.spawn_tracked(task_name, async move {
match local_data_fetcher.delete_directory(path.clone()).await {
Ok(_) => {
let _ = local_deleted_tx.send(None);
}
Err(e) => {
tracing::error!("Failed to delete local directory: {}", e);
let _ =
local_deleted_tx.send(Some(LocalError::from_message(e.to_string())));
}
}
}).await;
} else {
self.task_registry.spawn_tracked(task_name, async move {
match local_data_fetcher.delete_file(path.clone()).await {
Ok(_) => {
let _ = local_deleted_tx.send(None);
}
Err(e) => {
tracing::error!("Failed to delete local file: {}", e);
let _ =
local_deleted_tx.send(Some(LocalError::from_message(e.to_string())));
}
}
}).await;
}
}
async fn delete_s3_data(
&self,
item: S3SelectedItem,
s3_data_fetcher: S3DataFetcher,
s3_delete_tx: UnboundedSender<Option<S3Error>>,
) {
let items_with_children = flatten_items(vec![item]);
for item in items_with_children {
if !item.is_directory {
let delete_tx = s3_delete_tx.clone();
let fetcher = s3_data_fetcher.clone();
let task_name = format!("delete-s3:{}", item.name);
self.task_registry.spawn_tracked(task_name, async move {
match fetcher
.delete_data(
item.is_bucket,
item.bucket.clone(),
item.path.clone().unwrap_or(item.name.clone()),
item.is_directory,
)
.await
{
Ok(_) => {
let _ = delete_tx.send(None);
}
Err(e) => {
tracing::error!("Failed to delete S3 data: {}", e);
let _ = delete_tx.send(Some(S3Error::from_message(e.to_string())));
}
}
}).await;
}
}
}
async fn create_bucket(
&self,
name: String,
s3_data_fetcher: S3DataFetcher,
create_bucket_tx: UnboundedSender<Option<S3Error>>,
) {
let task_name = format!("create-bucket:{}", name);
self.task_registry.spawn_tracked(task_name, async move {
match s3_data_fetcher
.create_bucket(name.clone(), s3_data_fetcher.default_region.clone())
.await
{
Ok(_) => {
let _ = create_bucket_tx.send(None);
}
Err(e) => {
tracing::error!("Failed to create S3 bucket: {}", e);
let _ = create_bucket_tx.send(Some(S3Error::from_message(e.to_string())));
}
}
}).await;
}
fn get_current_s3_fetcher(
state: &State,
transfer_state_store: &Option<Arc<TransferStateStore>>,
) -> S3DataFetcher {
let fetcher = S3DataFetcher::new(state.current_creds.clone());
if let Some(store) = transfer_state_store {
fetcher.with_transfer_state_store(store.clone())
} else {
fetcher
}
}
pub async fn main_loop(
mut self,
mut terminator: Terminator,
mut action_rx: UnboundedReceiver<Action>,
mut interrupt_rx: broadcast::Receiver<Interrupted>,
creds: Vec<FileCredential>,
) -> eyre::Result<Interrupted> {
let local_data_fetcher = LocalDataFetcher::new();
let transfer_manager = Arc::new(TransferManager::new(S3_OPERATIONS_CONCURRENCY_LEVEL));
let job_item_map: Arc<Mutex<HashMap<JobId, TransferJobInfo>>> =
Arc::new(Mutex::new(HashMap::new()));
let transfer_state_store: Option<Arc<TransferStateStore>> =
match TransferStateStore::new(get_data_dir()).await {
Ok(store) => {
let (uploads, downloads) = store.get_pending_counts().await;
if uploads > 0 || downloads > 0 {
tracing::info!(
"Found {} pending uploads and {} pending downloads from previous session",
uploads,
downloads
);
}
Some(Arc::new(store))
}
Err(e) => {
tracing::warn!("Failed to initialize transfer state store: {}", e);
None
}
};
self.transfer_state_store = transfer_state_store.clone();
let mut state = State::new(creds.clone());
let s3_data_fetcher = Self::get_current_s3_fetcher(&state, &transfer_state_store);
state.set_s3_loading(true);
state.set_current_local_path(
dirs::home_dir()
.unwrap()
.as_path()
.to_string_lossy()
.to_string(),
);
let transfer_persistence = Arc::new(TransferPersistence::new(get_data_dir()));
if let Ok(persisted) = transfer_persistence.load().await {
if !persisted.s3_selected_items.is_empty() || !persisted.local_selected_items.is_empty() {
tracing::info!(
"Restoring {} downloads and {} uploads from previous session",
persisted.s3_selected_items.len(),
persisted.local_selected_items.len()
);
state.s3_selected_items = persisted.s3_selected_items;
state.local_selected_items = persisted.local_selected_items;
}
}
let (s3_tx, mut s3_rx) =
mpsc::unbounded_channel::<(Option<String>, Option<String>, Vec<S3DataItem>)>();
let (s3_full_list_tx, mut s3_full_list_rx) =
mpsc::unbounded_channel::<(Option<String>, Option<String>, Vec<S3DataItem>)>();
let (s3_deleted_tx, mut s3_deleted_rx) = mpsc::unbounded_channel::<Option<S3Error>>();
let (local_tx, mut local_rx) = mpsc::unbounded_channel::<(String, Vec<LocalDataItem>)>();
let (local_deleted_tx, mut local_deleted_rx) =
mpsc::unbounded_channel::<Option<LocalError>>();
let (upload_tx, mut upload_rx) = mpsc::channel::<UploadProgressItem>(PROGRESS_CHANNEL_CAPACITY);
let (download_tx, mut download_rx) = mpsc::channel::<DownloadProgressItem>(PROGRESS_CHANNEL_CAPACITY);
let (create_bucket_tx, mut create_bucket_rx) = mpsc::unbounded_channel::<Option<S3Error>>();
let (transfer_result_tx, mut transfer_result_rx) =
mpsc::unbounded_channel::<TransferResult>();
Self::spawn_transfer_worker(
&self.task_registry,
transfer_manager.clone(),
s3_data_fetcher.clone(),
job_item_map.clone(),
transfer_result_tx,
upload_tx.clone(),
download_tx.clone(),
)
.await;
self.fetch_s3_data(None, None, s3_data_fetcher.clone(), s3_tx.clone())
.await;
self.fetch_local_data(
Some(
dirs::home_dir()
.unwrap()
.as_path()
.to_string_lossy()
.to_string(),
),
local_data_fetcher.clone(),
local_tx.clone(),
)
.await;
self.state_tx.send(state.clone()).await?;
let _ticker = tokio::time::interval(Duration::from_secs(1));
let result = loop {
tokio::select! {
Some(action) = action_rx.recv() => match action {
Action::Exit => {
if let Err(e) = transfer_persistence.save(
&state.s3_selected_items,
&state.local_selected_items
).await {
tracing::warn!("Failed to save pending transfers: {}", e);
}
let _ = terminator.terminate(Interrupted::UserInt);
break Interrupted::UserInt;
},
Action::Navigate { page} => {
state.set_active_page(page.clone());
let _ = self.state_tx.send(state.clone()).await;
if page == ActivePage::FileManager {
state.set_s3_loading(true);
let _ = self.state_tx.send(state.clone()).await;
let s3_data_fetcher = Self::get_current_s3_fetcher(&state, &transfer_state_store);
self.fetch_s3_data(
state.current_s3_bucket.clone(),
state.current_s3_path.clone(),
s3_data_fetcher,
s3_tx.clone()
).await;
}
}
Action::FetchLocalData { path} =>
self.fetch_local_data(Some(path), local_data_fetcher.clone(), local_tx.clone()).await,
Action::FetchS3Data { bucket, prefix } => {
state.set_s3_loading(true);
let _ = self.state_tx.send(state.clone()).await;
let s3_data_fetcher = Self::get_current_s3_fetcher(&state, &transfer_state_store);
self.fetch_s3_data(bucket, prefix, s3_data_fetcher, s3_tx.clone()).await
}
Action::ListS3DataRecursiveForItem { item } => {
state.set_s3_list_recursive_loading(true);
let _ = self.state_tx.send(state.clone()).await;
let s3_data_fetcher = Self::get_current_s3_fetcher(&state, &transfer_state_store);
self.list_s3_data_recursive(item, s3_data_fetcher, s3_full_list_tx.clone()).await
}
Action::MoveBackLocal => self.move_back_local_data(state.current_local_path.clone(), local_data_fetcher.clone(), local_tx.clone()).await,
Action::SelectS3Item { item} => {
state.add_s3_selected_item(item);
let _ = self.state_tx.send(state.clone()).await;
},
Action::UnselectS3Item { item} => {
state.remove_s3_selected_item(item);
let _ = self.state_tx.send(state.clone()).await;
},
Action::SelectLocalItem { item} => {
state.add_local_selected_item(item);
let _ = self.state_tx.send(state.clone()).await;
},
Action::UnselectLocalItem { item } => {
state.remove_local_selected_item(item);
let _ = self.state_tx.send(state.clone()).await;
},
Action::RunTransfers => {
state.remove_already_transferred_items();
let updated_s3_items = Self::enqueue_downloads(
&transfer_manager,
state.s3_selected_items.clone(),
&job_item_map,
).await;
for updated_item in updated_s3_items {
state.update_s3_item_job_id(&updated_item);
}
let updated_local_items = Self::enqueue_uploads(
&transfer_manager,
state.local_selected_items.clone(),
&job_item_map,
).await;
for updated_item in updated_local_items {
state.update_local_item_job_id(&updated_item);
}
self.state_tx.send(state.clone()).await?;
let _ = transfer_persistence.save(
&state.s3_selected_items,
&state.local_selected_items
).await;
},
Action::SelectCurrentS3Creds { item} => {
state.set_current_s3_creds(item);
let _ = self.state_tx.send(state.clone()).await;
let s3_data_fetcher = Self::get_current_s3_fetcher(&state, &transfer_state_store);
self.fetch_s3_data(None, None, s3_data_fetcher, s3_tx.clone()).await;
},
Action::DeleteS3Item { item} => {
let s3_data_fetcher = Self::get_current_s3_fetcher(&state, &transfer_state_store);
self.delete_s3_data(item.clone(), s3_data_fetcher.clone(), s3_deleted_tx.clone()).await;
if item.is_bucket {
self.fetch_s3_data(None, None, s3_data_fetcher, s3_tx.clone()).await;
} else {
self.fetch_s3_data(item.bucket, None, s3_data_fetcher, s3_tx.clone()).await;
}
},
Action::DeleteLocalItem {item} => {
state.remove_local_selected_item(item.clone());
let _ = self.state_tx.send(state.clone()).await;
self.delete_local_data(item.clone(), local_data_fetcher.clone(), local_deleted_tx.clone()).await;
self.fetch_local_data(Some(item.path.clone()), local_data_fetcher.clone(), local_tx.clone()).await;
},
Action::CreateBucket {name} => {
let s3_data_fetcher = Self::get_current_s3_fetcher(&state, &transfer_state_store);
self.create_bucket(name.clone(), s3_data_fetcher.clone(), create_bucket_tx.clone()).await;
self.fetch_s3_data(None, None, s3_data_fetcher, s3_tx.clone()).await;
},
Action::ClearDeletionErrors => {
state.s3_delete_error = None;
state.local_delete_error = None;
state.create_bucket_error = None;
self.state_tx.send(state.clone()).await?;
},
Action::SortS3 { column } => {
state.sort_s3_data(column);
self.state_tx.send(state.clone()).await?;
},
Action::SortLocal { column } => {
state.sort_local_data(column);
self.state_tx.send(state.clone()).await?;
},
Action::SetSearchMode { active } => {
state.set_search_mode(active);
self.state_tx.send(state.clone()).await?;
},
Action::SetSearchQuery { query } => {
state.set_search_query(query);
self.state_tx.send(state.clone()).await?;
},
Action::ClearSearch => {
state.clear_search();
self.state_tx.send(state.clone()).await?;
},
Action::PauseTransfer { job_id } => {
if let Err(e) = transfer_manager.pause(job_id).await {
tracing::warn!("Failed to pause transfer {}: {}", job_id, e);
} else {
state.set_transfer_paused(job_id);
self.state_tx.send(state.clone()).await?;
let _ = transfer_persistence.save(
&state.s3_selected_items,
&state.local_selected_items
).await;
}
},
Action::ResumeTransfer { job_id } => {
if let Err(e) = transfer_manager.resume(job_id).await {
tracing::warn!("Failed to resume transfer {}: {}", job_id, e);
} else {
state.set_transfer_resumed(job_id);
self.state_tx.send(state.clone()).await?;
}
},
Action::CancelTransfer { job_id } => {
if let Err(e) = transfer_manager.cancel(job_id).await {
tracing::warn!("Failed to cancel transfer {}: {}", job_id, e);
} else {
state.set_transfer_cancelled(job_id);
job_item_map.lock().await.remove(&job_id);
self.state_tx.send(state.clone()).await?;
let _ = transfer_persistence.save(
&state.s3_selected_items,
&state.local_selected_items
).await;
}
},
},
Some(result) = transfer_result_rx.recv() => {
match result {
TransferResult::UploadComplete(item) => {
let dest_bucket = item.destination_bucket.clone();
state.update_selected_local_transfers(item);
if state.all_uploads_complete_for_bucket(&dest_bucket) {
let s3_data_fetcher = Self::get_current_s3_fetcher(&state, &transfer_state_store);
self.fetch_s3_data(Some(dest_bucket), None, s3_data_fetcher, s3_tx.clone()).await;
}
self.state_tx.send(state.clone()).await?;
let _ = transfer_persistence.save(
&state.s3_selected_items,
&state.local_selected_items
).await;
},
TransferResult::UploadFailed(item, _error) => {
state.update_selected_local_transfers(item);
self.state_tx.send(state.clone()).await?;
let _ = transfer_persistence.save(
&state.s3_selected_items,
&state.local_selected_items
).await;
},
TransferResult::DownloadComplete(item) => {
let dest_dir = item.destination_dir.clone();
state.update_selected_s3_transfers(item);
if state.all_downloads_complete_for_directory(&dest_dir) {
self.fetch_local_data(Some(dest_dir), local_data_fetcher.clone(), local_tx.clone()).await;
}
self.state_tx.send(state.clone()).await?;
let _ = transfer_persistence.save(
&state.s3_selected_items,
&state.local_selected_items
).await;
},
TransferResult::DownloadFailed(item, _error) => {
state.update_selected_s3_transfers(item);
self.state_tx.send(state.clone()).await?;
let _ = transfer_persistence.save(
&state.s3_selected_items,
&state.local_selected_items
).await;
},
}
},
Some((bucket, prefix, data)) = s3_rx.recv() => {
state.update_buckets(bucket, prefix, data);
self.state_tx.send(state.clone()).await?;
},
Some((_bucket, _prefix, data)) = s3_full_list_rx.recv() => {
state.update_s3_recursive_list(data);
self.state_tx.send(state.clone()).await?;
},
Some((path, files)) = local_rx.recv() => {
state.update_files(path, files);
self.state_tx.send(state.clone()).await?;
},
Some(item) = upload_rx.recv() => {
if state.active_page == ActivePage::Transfers {
state.update_progress_on_selected_local_item(&item);
self.state_tx.send(state.clone()).await?;
}
},
Some(item) = download_rx.recv() => {
if state.active_page == ActivePage::Transfers {
state.update_progress_on_selected_s3_item(&item);
self.state_tx.send(state.clone()).await?;
}
},
Some(error) = local_deleted_rx.recv() => {
state.set_local_delete_error(error);
self.state_tx.send(state.clone()).await?;
},
Some(error) = s3_deleted_rx.recv() => {
state.set_s3_delete_error(error);
self.state_tx.send(state.clone()).await?;
},
Some(error) = create_bucket_rx.recv() => {
state.set_create_bucket_error(error);
self.state_tx.send(state.clone()).await?;
}
Ok(interrupted) = interrupt_rx.recv() => {
break interrupted;
}
}
};
Ok(result)
}
}
#[derive(Debug, Clone)]
enum TransferJobInfo {
Upload { item: LocalSelectedItem },
Download { item: S3SelectedItem },
}