local: Hack to hopefully fix too many fds
This commit is contained in:
parent
6356d96e9c
commit
369d0135c3
59
src/local.rs
59
src/local.rs
|
@ -17,12 +17,49 @@ struct InfoPayload {
|
||||||
app_version: u32,
|
app_version: u32,
|
||||||
}
|
}
|
||||||
|
|
||||||
|
macro_rules! multi_error {
|
||||||
|
($name:tt { $($variant:tt($t:ty)),+ $(,)? }) => {
|
||||||
|
#[derive(Debug)]
|
||||||
|
pub enum $name {
|
||||||
|
$($variant($t)),+
|
||||||
|
}
|
||||||
|
|
||||||
|
impl ::std::fmt::Display for $name {
|
||||||
|
fn fmt(&self, f: &mut ::std::fmt::Formatter<'_>) -> ::std::fmt::Result {
|
||||||
|
match self {
|
||||||
|
$(
|
||||||
|
Self::$variant(x) => x.fmt(f)
|
||||||
|
),+
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
impl ::std::error::Error for $name {}
|
||||||
|
|
||||||
|
$(
|
||||||
|
impl From<$t> for $name {
|
||||||
|
fn from(value: $t) -> Self {
|
||||||
|
Self::$variant(value)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
)+
|
||||||
|
};
|
||||||
|
}
|
||||||
|
|
||||||
|
multi_error!(
|
||||||
|
Error {
|
||||||
|
Io(std::io::Error),
|
||||||
|
Http(reqwest::Error),
|
||||||
|
}
|
||||||
|
);
|
||||||
|
|
||||||
|
/// Represents a connection to the server running on the mobile device.
|
||||||
pub struct LocalServer {
|
pub struct LocalServer {
|
||||||
base_uri: reqwest::Url,
|
base_uri: reqwest::Url,
|
||||||
info: InfoPayload,
|
info: InfoPayload,
|
||||||
client: reqwest::Client,
|
client: reqwest::Client,
|
||||||
cache: FileCache,
|
cache: FileCache,
|
||||||
tasks: JoinSet<reqwest::Result<()>>,
|
tasks: JoinSet<Result<(), Error>>,
|
||||||
semaphore: Arc<Semaphore>,
|
semaphore: Arc<Semaphore>,
|
||||||
}
|
}
|
||||||
|
|
||||||
|
@ -72,10 +109,6 @@ impl LocalServer {
|
||||||
path: impl AsRef<Path>,
|
path: impl AsRef<Path>,
|
||||||
) -> Result<(), Box<dyn std::error::Error>> {
|
) -> Result<(), Box<dyn std::error::Error>> {
|
||||||
let path = path.as_ref();
|
let path = path.as_ref();
|
||||||
let fd = tokio::fs::OpenOptions::new()
|
|
||||||
.read(true)
|
|
||||||
.open(path.to_owned())
|
|
||||||
.await?;
|
|
||||||
if !self.should_upload(path) {
|
if !self.should_upload(path) {
|
||||||
return Err(Box::new(std::io::Error::new(
|
return Err(Box::new(std::io::Error::new(
|
||||||
std::io::ErrorKind::InvalidInput,
|
std::io::ErrorKind::InvalidInput,
|
||||||
|
@ -83,10 +116,6 @@ impl LocalServer {
|
||||||
)));
|
)));
|
||||||
}
|
}
|
||||||
let filename = path.file_name().unwrap().to_string_lossy().to_string();
|
let filename = path.file_name().unwrap().to_string_lossy().to_string();
|
||||||
|
|
||||||
let form = multipart::Form::new()
|
|
||||||
.part("filename", multipart::Part::text(filename.to_string()))
|
|
||||||
.part("file", multipart::Part::stream(fd).file_name(filename));
|
|
||||||
let client = self.client.clone();
|
let client = self.client.clone();
|
||||||
let base_uri = self.base_uri.clone();
|
let base_uri = self.base_uri.clone();
|
||||||
let semaphore = self.semaphore.clone();
|
let semaphore = self.semaphore.clone();
|
||||||
|
@ -96,6 +125,14 @@ impl LocalServer {
|
||||||
// The actual upload task
|
// The actual upload task
|
||||||
self.tasks.spawn(async move {
|
self.tasks.spawn(async move {
|
||||||
let _permit = semaphore.acquire_owned().await.unwrap();
|
let _permit = semaphore.acquire_owned().await.unwrap();
|
||||||
|
let fd = tokio::fs::OpenOptions::new()
|
||||||
|
.read(true)
|
||||||
|
.open(path.to_owned())
|
||||||
|
.await?;
|
||||||
|
|
||||||
|
let form = multipart::Form::new()
|
||||||
|
.part("filename", multipart::Part::text(filename.to_string()))
|
||||||
|
.part("file", multipart::Part::stream(fd).file_name(filename));
|
||||||
let response = client
|
let response = client
|
||||||
.post(base_uri.join("upload").unwrap())
|
.post(base_uri.join("upload").unwrap())
|
||||||
.multipart(form)
|
.multipart(form)
|
||||||
|
@ -110,7 +147,7 @@ impl LocalServer {
|
||||||
info!("Uploaded {}.", path.display());
|
info!("Uploaded {}.", path.display());
|
||||||
}
|
}
|
||||||
|
|
||||||
reqwest::Result::Ok(())
|
Ok(())
|
||||||
});
|
});
|
||||||
|
|
||||||
Ok(())
|
Ok(())
|
||||||
|
@ -118,7 +155,7 @@ impl LocalServer {
|
||||||
|
|
||||||
/// Provides an awaitable function for the task queue. This function will
|
/// Provides an awaitable function for the task queue. This function will
|
||||||
/// only return once the queue is empty or if an error occurs.
|
/// only return once the queue is empty or if an error occurs.
|
||||||
pub async fn wait_on_queue(&mut self) -> reqwest::Result<()> {
|
pub async fn wait_on_queue(&mut self) -> Result<(), Error> {
|
||||||
while let Some(task) = self.tasks.join_next().await {
|
while let Some(task) = self.tasks.join_next().await {
|
||||||
task.expect("upload task spawn error")?;
|
task.expect("upload task spawn error")?;
|
||||||
}
|
}
|
||||||
|
|
Loading…
Reference in a new issue