mirror of
https://github.com/Drop-OSS/drop.git
synced 2026-07-27 02:04:39 +10:00
feat: formatting
This commit is contained in:
+4
-5
@@ -2,18 +2,17 @@ use std::fs::{self, read_dir};
|
|||||||
|
|
||||||
use protobuf_codegen::Codegen;
|
use protobuf_codegen::Codegen;
|
||||||
|
|
||||||
const OUT_DIR: &'static str = "./src/proto/";
|
const OUT_DIR: &str = "./src/proto/";
|
||||||
|
|
||||||
fn main() {
|
fn main() {
|
||||||
let files = read_dir("./proto").unwrap();
|
let files = read_dir("./proto").unwrap();
|
||||||
let files = files.map(|v| format!("proto/{}", v.unwrap().file_name().into_string().unwrap()));
|
let files = files.map(|v| format!("proto/{}", v.unwrap().file_name().into_string().unwrap()));
|
||||||
|
|
||||||
read_dir(OUT_DIR).unwrap().into_iter().for_each(|v| {
|
read_dir(OUT_DIR).unwrap().for_each(|v| {
|
||||||
if let Ok(entry) = v {
|
if let Ok(entry) = v
|
||||||
if entry.file_name().to_str().unwrap().ends_with(".rs") {
|
&& entry.file_name().to_str().unwrap().ends_with(".rs") {
|
||||||
fs::remove_file(entry.path()).unwrap();
|
fs::remove_file(entry.path()).unwrap();
|
||||||
}
|
}
|
||||||
}
|
|
||||||
});
|
});
|
||||||
|
|
||||||
Codegen::new()
|
Codegen::new()
|
||||||
|
|||||||
@@ -22,11 +22,12 @@ pub struct DownloadContext {
|
|||||||
last_access: Instant,
|
last_access: Instant,
|
||||||
}
|
}
|
||||||
impl DownloadContext {
|
impl DownloadContext {
|
||||||
|
#[must_use]
|
||||||
pub fn last_access(&self) -> Instant {
|
pub fn last_access(&self) -> Instant {
|
||||||
self.last_access
|
self.last_access
|
||||||
}
|
}
|
||||||
pub fn reset_last_access(&mut self) {
|
pub fn reset_last_access(&mut self) {
|
||||||
self.last_access = Instant::now()
|
self.last_access = Instant::now();
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -71,7 +72,7 @@ fn create_backend(
|
|||||||
create_backend_constructor(&version_path).ok_or(StatusCode::INTERNAL_SERVER_ERROR)?;
|
create_backend_constructor(&version_path).ok_or(StatusCode::INTERNAL_SERVER_ERROR)?;
|
||||||
|
|
||||||
let backend = backend()
|
let backend = backend()
|
||||||
.inspect_err(|err| warn!("failed to create version backend: {:?}", err))
|
.inspect_err(|err| warn!("failed to create version backend: {err:?}"))
|
||||||
.map_err(|_| StatusCode::INTERNAL_SERVER_ERROR)?;
|
.map_err(|_| StatusCode::INTERNAL_SERVER_ERROR)?;
|
||||||
Ok(backend)
|
Ok(backend)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -67,9 +67,9 @@ impl AsyncRead for SpeedtestStream {
|
|||||||
match amount {
|
match amount {
|
||||||
Ok(amount) => self.remaining -= amount,
|
Ok(amount) => self.remaining -= amount,
|
||||||
Err(err) => return Poll::Ready(Err(err)),
|
Err(err) => return Poll::Ready(Err(err)),
|
||||||
};
|
}
|
||||||
};
|
}
|
||||||
return Poll::Ready(Ok(()));
|
Poll::Ready(Ok(()))
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -49,7 +49,7 @@ impl<'a, T: Stream> SemaphoreStream<'a, T> {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
impl<'a, T: Stream> Stream for SemaphoreStream<'a, T>
|
impl<T: Stream> Stream for SemaphoreStream<'_, T>
|
||||||
where
|
where
|
||||||
T: Stream,
|
T: Stream,
|
||||||
{
|
{
|
||||||
@@ -160,7 +160,7 @@ async fn get_file_reader(
|
|||||||
)
|
)
|
||||||
.await
|
.await
|
||||||
.map_err(|v| {
|
.map_err(|v| {
|
||||||
error!("reader error for '{}': {v:?}", relative_filename);
|
error!("reader error for '{relative_filename}': {v:?}");
|
||||||
StatusCode::INTERNAL_SERVER_ERROR
|
StatusCode::INTERNAL_SERVER_ERROR
|
||||||
})
|
})
|
||||||
}
|
}
|
||||||
@@ -181,13 +181,13 @@ async fn get_or_create_context<'a>(
|
|||||||
if let Some(already_done) = context_cache.get_mut(&key) {
|
if let Some(already_done) = context_cache.get_mut(&key) {
|
||||||
Ok(already_done)
|
Ok(already_done)
|
||||||
} else {
|
} else {
|
||||||
info!("generating context for {}...", game_id);
|
info!("generating context for {game_id}...");
|
||||||
let context_result =
|
let context_result =
|
||||||
create_download_context(state, game_id.clone(), version_name.clone()).await?;
|
create_download_context(state, game_id.clone(), version_name.clone()).await?;
|
||||||
|
|
||||||
state.context_cache.insert(key.clone(), context_result);
|
state.context_cache.insert(key.clone(), context_result);
|
||||||
|
|
||||||
info!("continuing download for {}", game_id);
|
info!("continuing download for {game_id}");
|
||||||
|
|
||||||
drop(permit);
|
drop(permit);
|
||||||
|
|
||||||
|
|||||||
@@ -1,10 +1,9 @@
|
|||||||
use std::{
|
use std::{
|
||||||
path::Path,
|
path::Path,
|
||||||
sync::{Arc, LazyLock},
|
sync::Arc,
|
||||||
};
|
};
|
||||||
|
|
||||||
use anyhow::anyhow;
|
use anyhow::anyhow;
|
||||||
use dashmap::DashMap;
|
|
||||||
use droplet_rs::versions::types::VersionBackend;
|
use droplet_rs::versions::types::VersionBackend;
|
||||||
use protobuf::Message;
|
use protobuf::Message;
|
||||||
|
|
||||||
@@ -48,7 +47,7 @@ pub async fn has_backend_rpc(
|
|||||||
|
|
||||||
fn create_backend(path: &String) -> Result<Box<dyn VersionBackend + Send + Sync>, anyhow::Error> {
|
fn create_backend(path: &String) -> Result<Box<dyn VersionBackend + Send + Sync>, anyhow::Error> {
|
||||||
let backend_constructor = droplet_rs::versions::create_backend_constructor(Path::new(path))
|
let backend_constructor = droplet_rs::versions::create_backend_constructor(Path::new(path))
|
||||||
.ok_or(anyhow!("backend doesn't exist at path {}", path))?;
|
.ok_or(anyhow!("backend doesn't exist at path {path}"))?;
|
||||||
let backend = backend_constructor()?;
|
let backend = backend_constructor()?;
|
||||||
|
|
||||||
Ok(backend)
|
Ok(backend)
|
||||||
|
|||||||
@@ -19,10 +19,9 @@ use crate::{
|
|||||||
static READER_SEMAPHORE: LazyLock<Semaphore> = LazyLock::new(|| {
|
static READER_SEMAPHORE: LazyLock<Semaphore> = LazyLock::new(|| {
|
||||||
let cores = std::env::var("READER_THREADS")
|
let cores = std::env::var("READER_THREADS")
|
||||||
.ok()
|
.ok()
|
||||||
.map(|v| str::parse::<usize>(&v).ok())
|
.and_then(|v| str::parse::<usize>(&v).ok())
|
||||||
.flatten()
|
|
||||||
.unwrap_or(num_cpus::get() / 2);
|
.unwrap_or(num_cpus::get() / 2);
|
||||||
info!("using {} import threads", cores);
|
info!("using {cores} import threads");
|
||||||
Semaphore::new(cores)
|
Semaphore::new(cores)
|
||||||
});
|
});
|
||||||
|
|
||||||
|
|||||||
@@ -1,6 +1,6 @@
|
|||||||
use std::sync::Arc;
|
use std::sync::Arc;
|
||||||
|
|
||||||
use log::{info, warn};
|
use log::warn;
|
||||||
|
|
||||||
use crate::{proto::{core::{DropBoundType, TorrentialBound}, droplet::RpcError}, server::DropServer};
|
use crate::{proto::{core::{DropBoundType, TorrentialBound}, droplet::RpcError}, server::DropServer};
|
||||||
|
|
||||||
@@ -15,7 +15,7 @@ where
|
|||||||
let message_id = message.message_id.clone();
|
let message_id = message.message_id.clone();
|
||||||
let result = rpc(server.clone(), message).await;
|
let result = rpc(server.clone(), message).await;
|
||||||
if let Err(err) = result {
|
if let Err(err) = result {
|
||||||
warn!("manifest generation failed with err: {:?}", err);
|
warn!("manifest generation failed with err: {err:?}");
|
||||||
let mut manifest_err = RpcError::new();
|
let mut manifest_err = RpcError::new();
|
||||||
manifest_err.error = err.to_string();
|
manifest_err.error = err.to_string();
|
||||||
let _ = server
|
let _ = server
|
||||||
|
|||||||
@@ -25,7 +25,7 @@ async fn main() {
|
|||||||
initialise_logger();
|
initialise_logger();
|
||||||
|
|
||||||
if let Ok(working_directory) = std::env::var("WORKING_DIRECTORY") {
|
if let Ok(working_directory) = std::env::var("WORKING_DIRECTORY") {
|
||||||
info!("moving to working directory {}", working_directory);
|
info!("moving to working directory {working_directory}");
|
||||||
set_current_dir(working_directory).expect("failed to change working directory");
|
set_current_dir(working_directory).expect("failed to change working directory");
|
||||||
}
|
}
|
||||||
|
|
||||||
@@ -38,7 +38,7 @@ async fn main() {
|
|||||||
|
|
||||||
let shared_state = Arc::new(AppState {
|
let shared_state = Arc::new(AppState {
|
||||||
context_cache: DashMap::new(),
|
context_cache: DashMap::new(),
|
||||||
server: server,
|
server,
|
||||||
});
|
});
|
||||||
|
|
||||||
let interval_shared_state = shared_state.clone();
|
let interval_shared_state = shared_state.clone();
|
||||||
@@ -62,7 +62,7 @@ async fn main() {
|
|||||||
};
|
};
|
||||||
if last_access.elapsed().as_secs() >= CONTEXT_TTL {
|
if last_access.elapsed().as_secs() >= CONTEXT_TTL {
|
||||||
shared_state.context_cache.remove(&key);
|
shared_state.context_cache.remove(&key);
|
||||||
info!("cleaned context: {:?}", key);
|
info!("cleaned context: {key:?}");
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
@@ -70,7 +70,7 @@ async fn main() {
|
|||||||
|
|
||||||
let app = setup_app(shared_state);
|
let app = setup_app(shared_state);
|
||||||
|
|
||||||
serve(app).await.unwrap();
|
serve(app).await.expect("failed to serve app");
|
||||||
}
|
}
|
||||||
|
|
||||||
fn setup_app(shared_state: Arc<AppState>) -> Router {
|
fn setup_app(shared_state: Arc<AppState>) -> Router {
|
||||||
@@ -87,7 +87,9 @@ fn setup_app(shared_state: Arc<AppState>) -> Router {
|
|||||||
}
|
}
|
||||||
|
|
||||||
async fn serve(app: Router) -> Result<(), std::io::Error> {
|
async fn serve(app: Router) -> Result<(), std::io::Error> {
|
||||||
let listener = tokio::net::TcpListener::bind("0.0.0.0:5000").await.unwrap();
|
let listener = tokio::net::TcpListener::bind("0.0.0.0:5000")
|
||||||
|
.await
|
||||||
|
.expect("failed to bind tcp server");
|
||||||
info!("started depot server");
|
info!("started depot server");
|
||||||
axum::serve(listener, app).await
|
axum::serve(listener, app).await
|
||||||
}
|
}
|
||||||
@@ -96,5 +98,5 @@ fn initialise_logger() {
|
|||||||
SimpleLogger::new()
|
SimpleLogger::new()
|
||||||
.with_level(log::LevelFilter::Info)
|
.with_level(log::LevelFilter::Info)
|
||||||
.init()
|
.init()
|
||||||
.unwrap();
|
.expect("failed to init logger");
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -33,5 +33,5 @@ pub async fn fetch_instance_games(app_state: &AppState) -> Result<Vec<SkeletonGa
|
|||||||
|
|
||||||
let response: ServerGamesResponse = app_state.server.wait_for_message_id(&message_id).await?;
|
let response: ServerGamesResponse = app_state.server.wait_for_message_id(&message_id).await?;
|
||||||
|
|
||||||
return Ok(response.games);
|
Ok(response.games)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -87,7 +87,7 @@ impl DropServer {
|
|||||||
}
|
}
|
||||||
|
|
||||||
/**
|
/**
|
||||||
Long-lived subroutine that never returns, runs the recieve_loop and reconnects
|
Long-lived subroutine that never returns, runs the `recieve_loop` and reconnects
|
||||||
as necessary
|
as necessary
|
||||||
*/
|
*/
|
||||||
async fn recieve_subroutine(myself: Arc<DropServer>, read_stream: OwnedReadHalf) -> ! {
|
async fn recieve_subroutine(myself: Arc<DropServer>, read_stream: OwnedReadHalf) -> ! {
|
||||||
@@ -95,7 +95,7 @@ impl DropServer {
|
|||||||
|
|
||||||
loop {
|
loop {
|
||||||
if let Err(err) = Self::recieve_loop(myself.clone(), &mut buffered_reader).await {
|
if let Err(err) = Self::recieve_loop(myself.clone(), &mut buffered_reader).await {
|
||||||
warn!("server disconnected with error: {:?}", err);
|
warn!("server disconnected with error: {err:?}");
|
||||||
|
|
||||||
let (drop_stream, _) = myself
|
let (drop_stream, _) = myself
|
||||||
.server
|
.server
|
||||||
@@ -130,14 +130,11 @@ impl DropServer {
|
|||||||
|
|
||||||
let message = message.value();
|
let message = message.value();
|
||||||
|
|
||||||
match message.type_.unwrap() {
|
if message.type_.unwrap() == crate::proto::core::TorrentialBoundType::ERROR {
|
||||||
crate::proto::core::TorrentialBoundType::ERROR => {
|
Err(anyhow!(String::from_utf8(message.data.clone()).unwrap()))
|
||||||
return Err(anyhow!(String::from_utf8(message.data.clone()).unwrap()));
|
} else {
|
||||||
}
|
let response = T::parse_from_bytes(&message.data)?;
|
||||||
_ => {
|
Ok(response)
|
||||||
let response = T::parse_from_bytes(&message.data)?;
|
|
||||||
return Ok(response);
|
|
||||||
}
|
|
||||||
}
|
}
|
||||||
}
|
}
|
||||||
|
|
||||||
|
|||||||
@@ -1,7 +1,6 @@
|
|||||||
use std::{collections::HashMap, path::PathBuf, sync::Arc};
|
use std::sync::Arc;
|
||||||
|
|
||||||
use dashmap::DashMap;
|
use dashmap::DashMap;
|
||||||
use tokio::sync::{OnceCell, Semaphore};
|
|
||||||
|
|
||||||
use crate::{DownloadContext, server::DropServer};
|
use crate::{DownloadContext, server::DropServer};
|
||||||
|
|
||||||
|
|||||||
Reference in New Issue
Block a user