Compare commits

..
12 Commits
Author SHA1 Message Date
Bluemangoo 8d4b6e79f3 fix export path 2026-05-04 15:26:22 +08:00
Bluemangoo cc4e6f0615 fix temp create&remove 2026-04-29 21:15:44 +08:00
Bluemangoo cdf617b5c3 fix task clear 2026-04-24 13:12:03 +08:00
Bluemangoo 46cf5cdf4b tcp proxy 2026-04-24 12:50:40 +08:00
Bluemangoo adc3afc7d5 fix ram 2026-04-24 12:00:03 +08:00
Bluemangoo de7b63d366 multi-server 2026-04-24 00:32:59 +08:00
Bluemangoo e019948017 dockerfile, modify it when using 2026-04-23 19:42:09 +08:00
Bluemangoo e00b25b6c7 fix tcp server wait 2026-04-23 19:00:05 +08:00
Bluemangoo 59c50d5e08 fix path sep 2026-04-23 18:31:54 +08:00
Bluemangoo b892c3b46f vibe config example 2026-04-23 16:08:41 +08:00
Bluemangoo 2f1c7eb9d8 fix tons of bugs 2026-04-17 18:56:56 +08:00
Bluemangoo d83d611ab8 dynamic config 2026-04-17 16:59:45 +08:00
26 changed files with 1521 additions and 211 deletions
+20
View File
@@ -0,0 +1,20 @@
# Git files
.git
.gitignore
# Build artifacts
target/
# Logs
logs/
*.log
# Data
data/
# IDE
.idea/
.vscode/
*.swp
*.swo
*~
Generated
+38
View File
@@ -122,9 +122,21 @@ dependencies = [
"tempfile", "tempfile",
"thiserror 2.0.18", "thiserror 2.0.18",
"tokio", "tokio",
"twox-hash",
"yaml_serde", "yaml_serde",
] ]
[[package]]
name = "async-http-proxy"
version = "1.2.5"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "29faa5d4d308266048bd7505ba55484315a890102f9345b9ff4b87de64201592"
dependencies = [
"httparse",
"thiserror 1.0.69",
"tokio",
]
[[package]] [[package]]
name = "async-lock" name = "async-lock"
version = "3.4.1" version = "3.4.1"
@@ -388,9 +400,11 @@ dependencies = [
"bytes", "bytes",
"common", "common",
"communicator", "communicator",
"futures-util",
"h2", "h2",
"lazy_static", "lazy_static",
"log", "log",
"reqwest",
"serde", "serde",
"serde_json", "serde_json",
"simplelog", "simplelog",
@@ -441,6 +455,7 @@ name = "communicator"
version = "0.1.0" version = "0.1.0"
dependencies = [ dependencies = [
"anyhow", "anyhow",
"async-http-proxy",
"bytes", "bytes",
"futures", "futures",
"futures-util", "futures-util",
@@ -452,7 +467,9 @@ dependencies = [
"serde_json", "serde_json",
"tokio", "tokio",
"tokio-rustls", "tokio-rustls",
"tokio-socks",
"tokio-util", "tokio-util",
"url",
"webpki-roots", "webpki-roots",
] ]
@@ -2712,6 +2729,18 @@ dependencies = [
"tokio", "tokio",
] ]
[[package]]
name = "tokio-socks"
version = "0.5.2"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "0d4770b8024672c1101b3f6733eab95b18007dbe0847a8afe341fcf79e06043f"
dependencies = [
"either",
"futures-util",
"thiserror 1.0.69",
"tokio",
]
[[package]] [[package]]
name = "tokio-util" name = "tokio-util"
version = "0.7.18" version = "0.7.18"
@@ -2795,6 +2824,15 @@ version = "0.2.5"
source = "registry+https://github.com/rust-lang/crates.io-index" source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "e421abadd41a4225275504ea4d6566923418b7f05506fbc9c0fe86ba7396114b" checksum = "e421abadd41a4225275504ea4d6566923418b7f05506fbc9c0fe86ba7396114b"
[[package]]
name = "twox-hash"
version = "2.1.2"
source = "registry+https://github.com/rust-lang/crates.io-index"
checksum = "9ea3136b675547379c4bd395ca6b938e5ad3c3d20fad76e7fe85f9e0d011419c"
dependencies = [
"rand",
]
[[package]] [[package]]
name = "typenum" name = "typenum"
version = "1.19.0" version = "1.19.0"
+5
View File
@@ -34,3 +34,8 @@ lazy_static = "1.5.0"
structopt = "0.3.26" structopt = "0.3.26"
tokio-util = "0.7.18" tokio-util = "0.7.18"
futures-util = "0.3.32" futures-util = "0.3.32"
twox-hash = "2.1.2"
futures = "0.3.32"
tokio-socks = "0.5.2"
url = "2.5.8"
async-http-proxy = "1.2.5"
+1
View File
@@ -22,3 +22,4 @@ cridecoder = { workspace = true }
thiserror = { workspace = true } thiserror = { workspace = true }
log = { workspace = true } log = { workspace = true }
cipher = { workspace = true, features = ["block-padding"] } cipher = { workspace = true, features = ["block-padding"] }
twox-hash = { workspace = true, features = ["xxhash3_128"] }
+11 -2
View File
@@ -262,13 +262,17 @@ impl AssetExecutionContext {
} }
RegionProviderConfig::Nuverse { RegionProviderConfig::Nuverse {
asset_version_url, asset_version_url,
app_version,
asset_info_url_template, asset_info_url_template,
.. ..
} => { } => {
// For nuverse, always fetch the version from asset_version_url. // For nuverse, always fetch the version from asset_version_url.
// The incoming request.asset_version is intentionally ignored here // The incoming request.asset_version is intentionally ignored here
// to match Go reference behavior. // to match Go reference behavior.
let app_version = self.sync_context.app_version.as_ref().ok_or_else(|| {
AssetExecutionError::MissingAppVersion {
region: self.region_name.clone(),
}
})?;
let version_url = asset_version_url.replace("{app_version}", app_version); let version_url = asset_version_url.replace("{app_version}", app_version);
let resolved_version = let resolved_version =
String::from_utf8_lossy(&self.get_with_retry(&version_url).await?) String::from_utf8_lossy(&self.get_with_retry(&version_url).await?)
@@ -319,9 +323,13 @@ impl AssetExecutionContext {
} }
RegionProviderConfig::Nuverse { RegionProviderConfig::Nuverse {
asset_bundle_url_template, asset_bundle_url_template,
app_version,
.. ..
} => { } => {
let app_version = self.sync_context.app_version.as_ref().ok_or_else(|| {
AssetExecutionError::MissingAppVersion {
region: self.region_name.clone(),
}
})?;
let asset_version = self let asset_version = self
.resolved_asset_version .resolved_asset_version
.as_deref() .as_deref()
@@ -435,6 +443,7 @@ impl AssetExecutionContext {
&self.region, &self.region,
&temp_file, &temp_file,
&task.bundle_path, &task.bundle_path,
&task.download_path,
category, category,
) )
.await; .await;
-1
View File
@@ -150,7 +150,6 @@ pub enum RegionProviderConfig {
}, },
Nuverse { Nuverse {
asset_version_url: String, asset_version_url: String,
app_version: String,
asset_info_url_template: String, asset_info_url_template: String,
asset_bundle_url_template: String, asset_bundle_url_template: String,
#[serde(default)] #[serde(default)]
+2
View File
@@ -107,6 +107,8 @@ pub enum AssetExecutionError {
HttpStatus { url: String, status: u16 }, HttpStatus { url: String, status: u16 },
#[error("region `{region}` is missing asset_save_dir")] #[error("region `{region}` is missing asset_save_dir")]
MissingAssetSaveDir { region: String }, MissingAssetSaveDir { region: String },
#[error("nuverse region `{region}` requires asset_version and asset_hash")]
MissingAppVersion { region: String },
#[error("colorful_palette region `{region}` requires asset_version and asset_hash")] #[error("colorful_palette region `{region}` requires asset_version and asset_hash")]
MissingAssetVersionOrHash { region: String }, MissingAssetVersionOrHash { region: String },
#[error("colorful_palette region `{region}` is missing profile hash for `{profile}`")] #[error("colorful_palette region `{region}` is missing profile hash for `{profile}`")]
+59 -64
View File
@@ -54,18 +54,39 @@ pub fn get_export_group(export_path: &str) -> &'static str {
"container" "container"
} }
pub fn get_hex_index(input: &str) -> String {
let hash_val = twox_hash::XxHash3_128::oneshot(input.as_bytes());
format!("{:032x}", hash_val)
}
pub fn empty_dir(base: PathBuf, name: String) -> PathBuf {
let mut dir = base.join(&name);
let mut cnt = 1;
while dir.exists() {
dir = base.join(format!("{}_{}", &name, cnt));
cnt += 1;
}
std::fs::create_dir_all(&dir).unwrap();
dir
}
pub async fn extract_unity_asset_bundle( pub async fn extract_unity_asset_bundle(
app_config: &AppConfig, app_config: &AppConfig,
sync_context: &SyncContext, sync_context: &SyncContext,
region: &RegionConfig, region: &RegionConfig,
asset_bundle_file: &Path, asset_bundle_file: &Path,
export_path: &str, export_path: &str,
download_path: &str,
category: &str, category: &str,
) -> Result<(PathBuf, bool), ExportPipelineError> { ) -> Result<(PathBuf, bool), ExportPipelineError> {
let output_dir = std::env::temp_dir() let hash = get_hex_index(export_path);
.join("sekai-updater") let output_dir = empty_dir(
.join("extract") std::env::temp_dir()
.join(&sync_context.region); .join("sekai-updater")
.join("extract")
.join(&sync_context.region),
hash,
);
let Some(asset_studio_cli_path) = app_config.tools.asset_studio_cli_path.as_deref() else { let Some(asset_studio_cli_path) = app_config.tools.asset_studio_cli_path.as_deref() else {
return Ok((asset_bundle_file.parent().unwrap().to_path_buf(), true)); return Ok((asset_bundle_file.parent().unwrap().to_path_buf(), true));
}; };
@@ -82,7 +103,7 @@ pub async fn extract_unity_asset_bundle(
}; };
let actual_export_path = if sync_context.export.by_category { let actual_export_path = if sync_context.export.by_category {
output_dir.join(category.to_lowercase()).join(export_path) output_dir.join(category.to_lowercase())
} else { } else {
output_dir.join(export_path) output_dir.join(export_path)
}; };
@@ -126,10 +147,23 @@ pub async fn extract_unity_asset_bundle(
.await?; .await?;
let mut count = 0; let mut count = 0;
walk(&actual_export_path, &mut |_| count += 1)?; let mut is_bundle_name_file = true;
let download_name = match download_path.rsplit_once("/") {
None => download_path,
Some((_, name)) => name,
};
walk(&output_dir, &mut |f| {
count += 1;
let file_name = f.file_name().unwrap().to_string_lossy();
is_bundle_name_file = is_bundle_name_file
&& match file_name.rsplit_once(".") {
Some((name, _)) => name.eq_ignore_ascii_case(download_name),
None => file_name.eq_ignore_ascii_case(download_name),
}
})?;
post_process_exported_files(app_config, sync_context, region, &actual_export_path).await?; post_process_exported_files(app_config, sync_context, region, &output_dir).await?;
Ok((actual_export_path, count <= 1)) Ok((output_dir, is_bundle_name_file && count <= 1))
} }
pub async fn post_process_exported_files( pub async fn post_process_exported_files(
@@ -145,7 +179,6 @@ pub async fn post_process_exported_files(
handle_usm_files( handle_usm_files(
export_path, export_path,
sync_context, sync_context,
region,
&app_config.tools.ffmpeg_path, &app_config.tools.ffmpeg_path,
&app_config.execution.retry, &app_config.execution.retry,
) )
@@ -168,7 +201,6 @@ pub async fn post_process_exported_files(
async fn handle_usm_files( async fn handle_usm_files(
export_path: &Path, export_path: &Path,
sync_context: &SyncContext, sync_context: &SyncContext,
region: &RegionConfig,
ffmpeg_path: &str, ffmpeg_path: &str,
retry: &crate::core::config::RetryConfig, retry: &crate::core::config::RetryConfig,
) -> Result<Vec<PathBuf>, ExportPipelineError> { ) -> Result<Vec<PathBuf>, ExportPipelineError> {
@@ -177,28 +209,17 @@ async fn handle_usm_files(
return Ok(Vec::new()); return Ok(Vec::new());
} }
let usm_input = if usm_files.len() == 1 { let mut out: Vec<PathBuf> = vec![];
usm_files[0].clone()
} else {
merge_usm_files(export_path, &usm_files)?
};
process_usm_file( for f in usm_files {
&usm_input, out.append(&mut process_usm_file(&f, sync_context, ffmpeg_path, retry).await?);
export_path, }
sync_context, Ok(out)
region,
ffmpeg_path,
retry,
)
.await
} }
async fn process_usm_file( async fn process_usm_file(
usm_file: &Path, usm_file: &Path,
export_path: &Path,
sync_context: &SyncContext, sync_context: &SyncContext,
_: &RegionConfig,
ffmpeg_path: &str, ffmpeg_path: &str,
retry: &crate::core::config::RetryConfig, retry: &crate::core::config::RetryConfig,
) -> Result<Vec<PathBuf>, ExportPipelineError> { ) -> Result<Vec<PathBuf>, ExportPipelineError> {
@@ -214,7 +235,10 @@ async fn process_usm_file(
if sync_context.export.video.convert_to_mp4 if sync_context.export.video.convert_to_mp4
&& sync_context.export.video.direct_usm_to_mp4_with_ffmpeg && sync_context.export.video.direct_usm_to_mp4_with_ffmpeg
{ {
let mp4 = export_path.join(format!("{output_name}.mp4")); let mp4 = usm_file
.parent()
.unwrap()
.join(format!("{output_name}.mp4"));
convert_usm_to_mp4(usm_file, &mp4, ffmpeg_path, retry).await?; convert_usm_to_mp4(usm_file, &mp4, ffmpeg_path, retry).await?;
remove_file_if_exists(usm_file)?; remove_file_if_exists(usm_file)?;
return Ok(vec![mp4]); return Ok(vec![mp4]);
@@ -226,7 +250,7 @@ async fn process_usm_file(
.and_then(|metadata| metadata.video_frame_rate()) .and_then(|metadata| metadata.video_frame_rate())
.filter(|(_, denominator)| *denominator > 0) .filter(|(_, denominator)| *denominator > 0)
.map(FrameRate::from_tuple); .map(FrameRate::from_tuple);
let extracted = codec::export_usm(usm_file, export_path)?; let extracted = codec::export_usm(usm_file, usm_file.parent().unwrap())?;
let mut generated = extracted.clone(); let mut generated = extracted.clone();
if sync_context.export.video.convert_to_mp4 { if sync_context.export.video.convert_to_mp4 {
@@ -237,7 +261,10 @@ async fn process_usm_file(
.map(|ext| ext.eq_ignore_ascii_case("m2v")) .map(|ext| ext.eq_ignore_ascii_case("m2v"))
.unwrap_or(false) .unwrap_or(false)
{ {
let mp4 = export_path.join(format!("{output_name}.mp4")); let mp4 = usm_file
.parent()
.unwrap()
.join(format!("{output_name}.mp4"));
convert_m2v_to_mp4( convert_m2v_to_mp4(
&extracted_file, &extracted_file,
&mp4, &mp4,
@@ -293,7 +320,7 @@ async fn handle_acb_files(
fn process_acb_file( fn process_acb_file(
acb_file: &Path, acb_file: &Path,
output_dir: &Path, _output_dir: &Path,
sync_context: &SyncContext, sync_context: &SyncContext,
region: &RegionConfig, region: &RegionConfig,
ffmpeg_path: &str, ffmpeg_path: &str,
@@ -347,7 +374,7 @@ fn process_acb_file(
&retry, &retry,
) )
})?; })?;
let final_outputs = move_result_files(output_dir, &generated)?; let final_outputs = move_result_files(acb_file.parent().unwrap(), &generated)?;
remove_file_if_exists(acb_file)?; remove_file_if_exists(acb_file)?;
Ok(final_outputs) Ok(final_outputs)
@@ -395,7 +422,7 @@ fn process_hca_file(
generated.clear(); generated.clear();
} }
let final_outputs = move_result_files(output_dir, &generated)?; let final_outputs = move_result_files(hca_file.parent().unwrap(), &generated)?;
Ok(final_outputs) Ok(final_outputs)
} }
@@ -676,38 +703,6 @@ fn panic_message(panic: Box<dyn std::any::Any + Send>) -> String {
} }
} }
fn merge_usm_files(dir: &Path, usm_files: &[PathBuf]) -> Result<PathBuf, ExportPipelineError> {
let dir_name = dir
.file_name()
.and_then(|name| name.to_str())
.unwrap_or("merged");
let merged_file = dir.join(format!("{dir_name}.usm"));
let mut target =
std::fs::File::create(&merged_file).map_err(|source| ExportPipelineError::Io {
path: merged_file.clone(),
source,
})?;
for source_path in usm_files {
if *source_path == merged_file {
continue;
}
let mut source =
std::fs::File::open(source_path).map_err(|source| ExportPipelineError::Io {
path: source_path.clone(),
source,
})?;
std::io::copy(&mut source, &mut target).map_err(|source| ExportPipelineError::Io {
path: source_path.clone(),
source,
})?;
drop(source);
remove_file_if_exists(source_path)?;
}
Ok(merged_file)
}
pub fn find_files(dir: &Path) -> Result<Vec<PathBuf>, ExportPipelineError> { pub fn find_files(dir: &Path) -> Result<Vec<PathBuf>, ExportPipelineError> {
let mut files = Vec::new(); let mut files = Vec::new();
walk(dir, &mut |path| { walk(dir, &mut |path| {
+20
View File
@@ -0,0 +1,20 @@
FROM rust:slim-bookworm AS builder
WORKDIR /usr/src/app
RUN apt-get update && apt-get install -y musl-tools pkg-config libssl-dev
RUN rustup target add x86_64-unknown-linux-musl
COPY . .
RUN cargo build --target x86_64-unknown-linux-musl --release --bin client
FROM alpine:latest
RUN apk add --no-cache ca-certificates
RUN update-ca-certificates
WORKDIR /app
COPY --from=builder /usr/src/app/target/x86_64-unknown-linux-musl/release/client .
# Copy config file
COPY sekai-unpacker-client.yaml .
CMD ["./client", "-p", "cn", "-p", "jp"]
EXPOSE 3000
+3 -1
View File
@@ -18,4 +18,6 @@ simplelog = { workspace = true }
anyhow = { workspace = true } anyhow = { workspace = true }
h2 = { workspace = true } h2 = { workspace = true }
serde_json = { workspace = true } serde_json = { workspace = true }
tokio-util = {workspace = true} tokio-util = { workspace = true }
reqwest = { workspace = true }
futures-util = { workspace = true }
+54
View File
@@ -11,6 +11,12 @@ pub struct ClientConfig {
pub profiles: HashMap<String, Profile>, pub profiles: HashMap<String, Profile>,
} }
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
pub struct DynamicLoad {
pub url: String,
pub map: HashMap<String, String>,
}
#[derive(Debug, Clone, Serialize, Deserialize)] #[derive(Debug, Clone, Serialize, Deserialize)]
pub struct Profile { pub struct Profile {
#[serde(flatten)] #[serde(flatten)]
@@ -18,4 +24,52 @@ pub struct Profile {
pub path: String, pub path: String,
pub interval: Option<u64>, pub interval: Option<u64>,
pub concurrent: Option<usize>, pub concurrent: Option<usize>,
dynamic_load: Option<DynamicLoad>,
}
impl Profile {
pub fn need_update(&self) -> bool {
self.dynamic_load.is_some()
}
pub async fn update(&mut self) -> anyhow::Result<()> {
let Some(load) = &self.dynamic_load else {
return Ok(());
};
let response = reqwest::get(&load.url).await?;
let body = response.text().await?;
let data: serde_json::Value = serde_json::from_str(&body)?;
if !data.is_object() {
return Err(anyhow::Error::msg("response is not an object"));
};
let obj = data.as_object().unwrap();
macro_rules! load_field {
(from $obj:expr; $($tail:tt)*) => {
load_field!(@munch $obj, $($tail)*);
};
(@munch $obj:expr, $field:ident$(.$f_more:ident)* = $key:ident as $t:ident $($rest:tt)*) => {
load_field!(@suffix $obj, ($field$(.$f_more)*), $key, $t, [], $($rest)*);
};
(@suffix $obj:expr, ($($f_full:tt)*), $key:ident, $t:ident, [$($acc:tt)*], ; $($rest:tt)*) => {
if let Some(k) = load.map.get(stringify!($key)) {
if let Some(serde_json::Value::$t(v)) = $obj.get(k) {
$($f_full)* = Some(v.clone() $($acc)*);
}
}
load_field!(@munch $obj, $($rest)*);
};
(@suffix $obj:expr, ($($f_full:tt)*), $key:ident, $t:ident, [$($acc:tt)*], $token:tt $($rest:tt)*) => {
load_field!(@suffix $obj, ($($f_full)*), $key, $t, [$($acc)* $token], $($rest)*);
};
(@munch $obj:expr, ) => {};
}
load_field! {
from obj;
self.sync_context.asset_version = asset_version as String;
self.sync_context.asset_hash = asset_hash as String;
self.sync_context.app_version = app_version as String;
self.concurrent = concurrent as Number.as_u64().ok_or(anyhow::Error::msg("concurrent is not usize"))? as usize;
}
Ok(())
}
} }
+145 -23
View File
@@ -1,23 +1,30 @@
use crate::config::{ClientConfig, Profile}; use crate::config::{ClientConfig, Profile};
use crate::task::run; use crate::queue::SharedQueue;
use crate::signal::Signal;
use crate::task::{AtomicCounters, AutoSaveManifest, pre_run, run_main, run_side};
use common::strings::REGION_NOT_FOUND; use common::strings::REGION_NOT_FOUND;
use common::updater::DownloadTask;
use communicator::{ClientManager, Identity, TunnelEndpoint, TunnelListener, connect_tunnel}; use communicator::{ClientManager, Identity, TunnelEndpoint, TunnelListener, connect_tunnel};
use futures_util::future::join_all;
use lazy_static::lazy_static; use lazy_static::lazy_static;
use log::{LevelFilter, error, info}; use log::{LevelFilter, error, info};
use simplelog::{ColorChoice, Config, TermLogger, TerminalMode}; use simplelog::{ColorChoice, Config, TermLogger, TerminalMode};
use std::collections::{HashMap, VecDeque}; use std::collections::{HashMap, VecDeque};
use std::fs; use std::fs;
use std::path::Path;
use std::str::FromStr; use std::str::FromStr;
use std::sync::{Arc, RwLock}; use std::sync::Arc;
use std::time::Duration; use std::time::Duration;
use structopt::StructOpt; use structopt::StructOpt;
use tokio::sync::{OwnedSemaphorePermit, Semaphore, mpsc}; use tokio::sync::{OwnedSemaphorePermit, RwLock, Semaphore, mpsc};
use tokio::task::JoinSet; use tokio::task::JoinSet;
use tokio::time::sleep; use tokio::time::sleep;
use tokio_util::sync::CancellationToken; use tokio_util::sync::CancellationToken;
mod config; mod config;
mod http; mod http;
mod queue;
mod signal;
mod task; mod task;
#[derive(StructOpt)] #[derive(StructOpt)]
@@ -51,17 +58,23 @@ async fn main() -> anyhow::Result<()> {
std::process::exit(1); std::process::exit(1);
} }
let profiles = command_opts let profiles = join_all(command_opts.profile.iter().map(async |p| {
.profile let profile = CONFIG.profiles.get(p.as_str());
.iter() match profile {
.map(|p| { Some(profile) => {
let profile = CONFIG.profiles.get(p.as_str()); let mut profile = profile.clone();
match profile { if profile.need_update() {
Some(profile) => Ok((p.to_string(), profile.clone())), info!("Loading async config of {}", p);
None => Err(anyhow::anyhow!("Profile `{}` not found in config", p)), profile.update().await?;
}
Ok((p.to_string(), Arc::new(RwLock::new(profile))))
} }
}) None => Err(anyhow::anyhow!("Profile `{}` not found in config", p)),
.collect::<Result<Vec<_>, _>>()?; }
}))
.await
.into_iter()
.collect::<Result<Vec<_>, _>>()?;
let mut join_set = JoinSet::new(); let mut join_set = JoinSet::new();
@@ -119,17 +132,39 @@ async fn main() -> anyhow::Result<()> {
for profile in profiles { for profile in profiles {
let profile = Arc::new(profile.clone()); let profile = Arc::new(profile.clone());
let semaphore = Arc::new(Semaphore::new(1)); let semaphore = Arc::new(Semaphore::new(1));
let tasks: SharedQueue<DownloadTask> = SharedQueue::new();
let cnt = AtomicCounters::new();
let local_manifest = Arc::new(
AutoSaveManifest::new(
5,
Path::new(&{ profile.1.read().await.path.clone() })
.join("manifest.json")
.to_path_buf(),
)
.await?,
);
let signal = Signal::new();
let cancel_token = CancellationToken::new(); let cancel_token = CancellationToken::new();
let post_task = { let post_task = {
async |profile: Arc<(String, Profile)>, async |profile: Arc<(String, Arc<RwLock<Profile>>)>,
permit: OwnedSemaphorePermit, permit: OwnedSemaphorePermit,
cancel_token: CancellationToken| { cancel_token: CancellationToken| {
match profile.1.interval { let (interval, need_update) = {
let p = profile.1.read().await;
(p.interval, p.need_update())
};
match interval {
None => { None => {
cancel_token.cancel(); cancel_token.cancel();
} }
Some(interval) => { Some(interval) => {
sleep(Duration::from_secs(interval)).await; sleep(Duration::from_secs(interval)).await;
if need_update {
info!("Trying to update profile for {}", profile.0);
if let Err(e) = profile.1.write().await.update().await {
error!("Failed to update profile for {}: {}", profile.0, e);
}
}
} }
} }
drop(permit); drop(permit);
@@ -141,16 +176,20 @@ async fn main() -> anyhow::Result<()> {
let semaphore = semaphore.clone(); let semaphore = semaphore.clone();
let cancel_token = cancel_token.clone(); let cancel_token = cancel_token.clone();
let profile = profile.clone(); let profile = profile.clone();
let signal = signal.clone();
let tasks = tasks.clone();
let cnt = cnt.clone();
let local_manifest = local_manifest.clone();
let liveness_tx = liveness_tx.clone(); let liveness_tx = liveness_tx.clone();
join_set.spawn(async move { join_set.spawn(async move {
let _guard = liveness_tx; let _guard = liveness_tx;
let mut inner_set = JoinSet::new();
loop { loop {
if cancel_token.is_cancelled() { if cancel_token.is_cancelled() {
return; return;
} }
let client = sender.recv(); let client = sender.recv();
if client.is_none() { if client.is_none() {
tokio::time::sleep(Duration::from_millis(50)).await;
continue; continue;
} }
let client = client.unwrap(); let client = client.unwrap();
@@ -160,7 +199,12 @@ async fn main() -> anyhow::Result<()> {
let semaphore = semaphore.clone(); let semaphore = semaphore.clone();
let cancel_token = cancel_token.clone(); let cancel_token = cancel_token.clone();
let profile = profile.clone(); let profile = profile.clone();
inner_set.spawn(async move { let tasks = tasks.clone();
let cnt = cnt.clone();
let local_manifest = local_manifest.clone();
let signal = signal.clone();
tokio::spawn(async move {
loop { loop {
if client.get_client().await.is_err() { if client.get_client().await.is_err() {
return; return;
@@ -169,11 +213,38 @@ async fn main() -> anyhow::Result<()> {
return; return;
} }
{ {
let permit = semaphore.clone().acquire_owned().await.unwrap(); let permit = loop {
tokio::select! {
awaitable = signal.subscribe() => {
let r = run_side(client.clone(), tasks.clone(),cnt.clone(),local_manifest.clone(),profile.clone()).await;
if let Err(e)=r{
error!("{}", e);
}
awaitable.wait().await;
}
res = semaphore.clone().acquire_owned() => {
let permit = res.expect("Semaphore closed");
break permit;
}
}
};
if cancel_token.is_cancelled() { if cancel_token.is_cancelled() {
return; return;
} }
let result = run(client.clone(), profile.clone()).await; let sync_id = pre_run(client.clone(), profile.clone(), tasks.clone(), cnt.clone()).await;
let sig = signal.pick().await;
let result = match sync_id {
Ok(Some(id)) => {
run_main(client.clone(), profile.clone(), id, tasks.clone(), cnt.clone(), local_manifest.clone()).await
}
Ok(None) => {
Ok(true)
}
Err(e) => { Err(e) }
};
drop(sig);
match result { match result {
Ok(true) => { Ok(true) => {
post_task(profile.clone(), permit, cancel_token.clone()) post_task(profile.clone(), permit, cancel_token.clone())
@@ -183,6 +254,16 @@ async fn main() -> anyhow::Result<()> {
if error.to_string().contains(REGION_NOT_FOUND) { if error.to_string().contains(REGION_NOT_FOUND) {
return; return;
} }
if error
.to_string()
.contains("Session did not reconnect within 15s")
{
error!(
"Session lost for profile {}. Waiting for a fresh tunnel endpoint...",
profile.0
);
return;
}
error!("{}", error); error!("{}", error);
} }
_ => {} _ => {}
@@ -200,6 +281,10 @@ async fn main() -> anyhow::Result<()> {
let semaphore = semaphore.clone(); let semaphore = semaphore.clone();
let cancel_token = cancel_token.clone(); let cancel_token = cancel_token.clone();
let profile = profile.clone(); let profile = profile.clone();
let tasks = tasks.clone();
let cnt = cnt.clone();
let local_manifest = local_manifest.clone();
let signal = signal.clone();
info!("tcp client started for {}", client_conf.url); info!("tcp client started for {}", client_conf.url);
join_set.spawn(async move { join_set.spawn(async move {
loop { loop {
@@ -223,17 +308,54 @@ async fn main() -> anyhow::Result<()> {
return; return;
} }
{ {
let permit = semaphore.clone().acquire_owned().await.unwrap(); let permit = loop {
tokio::select! {
awaitable = signal.subscribe() => {
let r = run_side(client.clone(), tasks.clone(),cnt.clone(),local_manifest.clone(),profile.clone()).await;
if let Err(e)=r{
error!("{}", e);
}
awaitable.wait().await;
}
res = semaphore.clone().acquire_owned() => {
let permit = res.expect("Semaphore closed");
break permit;
}
}
};
if cancel_token.is_cancelled() { if cancel_token.is_cancelled() {
return; return;
} }
let result = run(client.clone(), profile.clone()).await; let sync_id = pre_run(client.clone(), profile.clone(), tasks.clone(), cnt.clone()).await;
let sig = signal.pick().await;
let result = match sync_id {
Ok(Some(id)) => {
run_main(client.clone(), profile.clone(), id, tasks.clone(), cnt.clone(), local_manifest.clone()).await
}
Ok(None) => {
Ok(true)
}
Err(e) => { Err(e) }
};
drop(sig);
match result { match result {
Ok(true) => { Ok(true) => {
post_task(profile.clone(), permit, cancel_token.clone()) post_task(profile.clone(), permit, cancel_token.clone())
.await; .await;
} }
Err(error) => { Err(error) => {
if error
.to_string()
.contains("Session did not reconnect within 15s")
{
error!(
"Session lost for profile {}. Reconnecting tunnel...",
profile.0
);
break;
}
error!("{}", error); error!("{}", error);
} }
_ => {} _ => {}
@@ -263,13 +385,13 @@ async fn main() -> anyhow::Result<()> {
} }
struct Sender<T> { struct Sender<T> {
inner: RwLock<VecDeque<T>>, inner: std::sync::RwLock<VecDeque<T>>,
} }
impl<T> Sender<T> { impl<T> Sender<T> {
pub fn new() -> Self { pub fn new() -> Self {
Self { Self {
inner: RwLock::new(VecDeque::new()), inner: std::sync::RwLock::new(VecDeque::new()),
} }
} }
+136
View File
@@ -0,0 +1,136 @@
use std::collections::VecDeque;
use std::sync::{Arc, Condvar, Mutex, atomic::{AtomicUsize, Ordering}};
pub struct SharedQueue<T> {
inner: Arc<QueueInner<T>>,
}
struct QueueInner<T> {
data: Mutex<VecDeque<T>>,
// 用于 pop 的阻塞
pop_cond: Condvar,
// 用于“全部消费完”的阻塞
done_cond: Condvar,
// 在途任务计数(队列中 + 正在处理中)
pending: AtomicUsize,
}
/// 任务守卫:当它被释放时,说明消费彻底结束
pub struct TaskGuard<T> {
pub item: T,
inner: Arc<QueueInner<T>>,
}
impl<T: Clone> Clone for TaskGuard<T> {
fn clone(&self) -> Self {
// 关键:每多出一个 Guard 副本,就意味着多了一个需要等待的“消费行为”
// 必须增加全局在途计数,否则会导致 pending 减成负数或提前归零
self.inner.pending.fetch_add(1, Ordering::SeqCst);
Self {
item: self.item.clone(),
inner: self.inner.clone(),
}
}
}
impl<T> Drop for TaskGuard<T> {
fn drop(&mut self) {
// 1. 任务完成,计数减一
let prev = self.inner.pending.fetch_sub(1, Ordering::SeqCst);
// 2. 如果减完后是 0,说明最后一项任务也处理完了
if prev == 1 {
let _lock = self.inner.data.lock().unwrap();
self.inner.done_cond.notify_all();
}
}
}
impl<T> SharedQueue<T> {
pub fn new() -> Self {
Self {
inner: Arc::new(QueueInner {
data: Mutex::new(VecDeque::new()),
pop_cond: Condvar::new(),
done_cond: Condvar::new(),
pending: AtomicUsize::new(0),
}),
}
}
// pub fn push(&self, item: T) {
// let mut queue = self.inner.data.lock().unwrap();
// // 增加在途计数
// self.inner.pending.fetch_add(1, Ordering::SeqCst);
// queue.push_back(item);
// self.inner.pop_cond.notify_one();
// }
pub fn push_all(&self, items: impl IntoIterator<Item = T>) {
let mut queue = self.inner.data.lock().unwrap();
let mut count = 0;
for item in items {
queue.push_back(item);
count += 1;
}
if count > 0 {
self.inner.pending.fetch_add(count, Ordering::SeqCst);
self.inner.pop_cond.notify_all();
}
}
// pub fn pop(&self) -> TaskGuard<T> {
// let mut queue = self.inner.data.lock().unwrap();
// while queue.is_empty() {
// queue = self.inner.pop_cond.wait(queue).unwrap();
// }
// let item = queue.pop_front().unwrap();
// TaskGuard {
// item,
// inner: self.inner.clone(),
// }
// }
pub fn try_pop(&self) -> Option<TaskGuard<T>> {
let mut queue = self.inner.data.lock().unwrap();
queue.pop_front().map(|item| {
TaskGuard {
item,
inner: self.inner.clone(),
}
})
}
pub fn clear(&self) {
let mut queue = self.inner.data.lock().unwrap();
let removed_count = queue.len();
if removed_count > 0 {
queue.clear();
let prev_pending = self.inner.pending.fetch_sub(removed_count, Ordering::SeqCst);
if prev_pending == removed_count {
self.inner.done_cond.notify_all();
}
}
}
/// 阻塞当前线程,直到所有在途任务(pending == 0)全部处理完
pub fn wait_until_all_consumed(&self) {
let mut _queue_lock = self.inner.data.lock().unwrap();
while self.inner.pending.load(Ordering::SeqCst) > 0 {
_queue_lock = self.inner.done_cond.wait(_queue_lock).unwrap();
}
}
}
impl<T> Clone for SharedQueue<T> {
fn clone(&self) -> Self {
Self { inner: Arc::clone(&self.inner) }
}
}
+86
View File
@@ -0,0 +1,86 @@
use std::sync::Arc;
use tokio::sync::{watch, Mutex, OwnedMutexGuard};
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum Stage {
Idle, // 空闲/等待发令
Processing, // Leader 干活中
}
pub struct Signal {
inner: Arc<Inner>,
}
struct Inner {
stage_tx: watch::Sender<Stage>,
pick_lock: Arc<Mutex<()>>,
}
impl Signal {
pub fn new() -> Self {
let (stage_tx, _) = watch::channel(Stage::Idle);
Self {
inner: Arc::new(Inner {
stage_tx,
pick_lock: Arc::new(Mutex::new(())),
}),
}
}
pub async fn pick(&self) -> LeaderHandler {
let lock_handle = self.inner.pick_lock.clone();
let _owned_guard = lock_handle.lock_owned().await;
// 切换到工作状态
let _ = self.inner.stage_tx.send(Stage::Processing);
LeaderHandler {
inner: self.inner.clone(),
_guard: _owned_guard,
}
}
pub async fn subscribe(&self) -> FollowerAwaitable {
let mut rx = self.inner.stage_tx.subscribe();
// 如果当前是 Idle,就挂起等待 Leader 变为 Processing
while *rx.borrow() != Stage::Processing {
if rx.changed().await.is_err() {
break;
}
}
FollowerAwaitable { rx }
}
}
pub struct LeaderHandler {
inner: Arc<Inner>,
_guard: OwnedMutexGuard<()>,
}
impl Drop for LeaderHandler {
fn drop(&mut self) {
// Leader 掉落,重置为 Idle,允许下一轮竞争
let _ = self.inner.stage_tx.send(Stage::Idle);
}
}
pub struct FollowerAwaitable {
rx: watch::Receiver<Stage>,
}
impl FollowerAwaitable {
pub async fn wait(mut self) {
// 等待状态变回 Idle (说明 Leader 掉落了)
while *self.rx.borrow() == Stage::Processing {
if self.rx.changed().await.is_err() {
break;
}
}
}
}
impl Clone for Signal {
fn clone(&self) -> Self {
Self { inner: self.inner.clone() }
}
}
+210 -35
View File
@@ -1,6 +1,9 @@
use crate::config::Profile; use crate::config::Profile;
use crate::http::{close, download, sync}; use crate::http::{close, download, sync};
use crate::queue::SharedQueue;
use anyhow::anyhow;
use common::http::{CloseRequest, DownloadRequest}; use common::http::{CloseRequest, DownloadRequest};
use common::updater::DownloadTask;
use communicator::ClientManager; use communicator::ClientManager;
use log::{error, info}; use log::{error, info};
use serde::{Deserialize, Serialize}; use serde::{Deserialize, Serialize};
@@ -10,23 +13,21 @@ use std::sync::Arc;
use std::sync::atomic::{AtomicUsize, Ordering}; use std::sync::atomic::{AtomicUsize, Ordering};
use tokio::sync::{RwLock, Semaphore}; use tokio::sync::{RwLock, Semaphore};
use tokio::task::JoinSet; use tokio::task::JoinSet;
use tokio_util::sync::CancellationToken;
pub async fn run( pub async fn pre_run(
client: Arc<ClientManager>, client: Arc<ClientManager>,
profile: Arc<(String, Profile)>, profile: Arc<(String, Arc<RwLock<Profile>>)>,
) -> anyhow::Result<bool> { queue: SharedQueue<DownloadTask>,
cnt: AtomicCounters,
) -> anyhow::Result<Option<String>> {
info!("[{}]: Starting sync", profile.0); info!("[{}]: Starting sync", profile.0);
tokio::fs::create_dir_all(&profile.1.path).await?; let p1 = Arc::new(profile.1.read().await.clone());
let sync_resp = sync(&mut client.get_client().await?, &profile.1).await?; tokio::fs::create_dir_all(&p1.path).await?;
let sync_resp = sync(&mut client.get_client().await?, &p1).await?;
let id = sync_resp.id; let id = sync_resp.id;
let local_manifest = Arc::new( let local_manifest = Arc::new(
AutoSaveManifest::new( AutoSaveManifest::new(5, Path::new(&p1.path).join("manifest.json").to_path_buf()).await?,
5,
Path::new(&profile.1.path)
.join("manifest.json")
.to_path_buf(),
)
.await?,
); );
let manifest_snapshot = { local_manifest.manifest.read().await.clone() }; let manifest_snapshot = { local_manifest.manifest.read().await.clone() };
let all_cnt = sync_resp.tasks.len(); let all_cnt = sync_resp.tasks.len();
@@ -51,59 +52,183 @@ pub async fn run(
info!("[{}]: No tasks to sync, skipping", profile.0); info!("[{}]: No tasks to sync, skipping", profile.0);
let req = CloseRequest { id: id.clone() }; let req = CloseRequest { id: id.clone() };
close(&mut client.get_client().await?, &req).await?; close(&mut client.get_client().await?, &req).await?;
return Ok(true); return Ok(None);
} }
let n = profile.1.concurrent.unwrap_or(5); cnt.reset();
queue.clear();
queue.push_all(tasks);
Ok(Some(id))
}
pub async fn run_main(
client: Arc<ClientManager>,
profile: Arc<(String, Arc<RwLock<Profile>>)>,
id: String,
queue: SharedQueue<DownloadTask>,
cnt: AtomicCounters,
manifest: Arc<AutoSaveManifest>,
) -> anyhow::Result<bool> {
info!("[{}]: Starting sync", profile.0);
let p1 = Arc::new(profile.1.read().await.clone());
let n = p1.concurrent.unwrap_or(5);
info!("[{}]: Start sync with {} thread", profile.0, n); info!("[{}]: Start sync with {} thread", profile.0, n);
let semaphore = Arc::new(Semaphore::new(n)); let semaphore = Arc::new(Semaphore::new(n));
let mut join_set = JoinSet::new(); let mut join_set = JoinSet::new();
for task in tasks { let cancel_token = CancellationToken::new();
while let Some(task) = queue.try_pop() {
if cancel_token.is_cancelled() {
break;
}
let permit = semaphore.clone().acquire_owned().await?; let permit = semaphore.clone().acquire_owned().await?;
let client = client.clone(); let client = client.clone();
let id = id.clone(); let id = id.clone();
let local_manifest = local_manifest.clone(); let local_manifest = manifest.clone();
let profile = profile.clone(); let p1 = p1.clone();
let cancel_token = cancel_token.clone();
join_set.spawn(async move { join_set.spawn(async move {
if cancel_token.is_cancelled() {
return Ok::<(), anyhow::Error>(());
}
let guard = task;
let task = &guard.item;
let req = DownloadRequest { let req = DownloadRequest {
id: id.clone(), id: id.clone(),
task: task.clone(), task: task.clone(),
}; };
let result = download(&mut client.get_client().await.unwrap(), &req, &profile.1).await; let mut conn = client.get_client().await?;
if let Err(e) = result let mut result = download(&mut conn, &req, &p1).await;
&& let Some(_) = e.downcast_ref::<h2::Error>() if let Err(e) = &result
&& e.downcast_ref::<h2::Error>().is_some()
{ {
download(&mut client.get_client().await.unwrap(), &req, &profile.1) let mut retry_conn = client.get_client().await?;
.await result = download(&mut retry_conn, &req, &p1).await;
.unwrap();
} }
if let Err(e) = &result
&& e.downcast_ref::<h2::Error>().is_some()
{
cancel_token.cancel();
}
result?;
local_manifest local_manifest
.add_bundle(task.bundle_path.clone(), task.bundle_hash.clone()) .add_bundle(task.bundle_path.clone(), task.bundle_hash.clone())
.await .await?;
.unwrap();
drop(permit); drop(permit);
Ok::<(), anyhow::Error>(())
}); });
} }
let mut succeed = 0;
let mut failed = 0;
while let Some(r) = join_set.join_next().await { while let Some(r) = join_set.join_next().await {
if let Err(e) = r { match r {
error!("{}", e); Ok(Ok(())) => cnt.inc_success(),
failed += 1; Ok(Err(e)) => {
} else { error!("{}", e);
succeed += 1; cnt.inc_failure()
}
Err(e) => {
error!("{}", e);
cnt.inc_failure()
}
} }
} }
local_manifest.save().await?; manifest.save().await?;
queue.wait_until_all_consumed();
info!( info!(
"[{}]: Sync finished with {} succeed, {} failed", "[{}]: Sync finished with {} succeed, {} failed",
profile.0, succeed, failed profile.0,
cnt.get_success(),
cnt.get_failure()
); );
let req = CloseRequest { id: id.clone() }; let req = CloseRequest { id: id.clone() };
close(&mut client.get_client().await?, &req).await?; close(&mut client.get_client().await?, &req).await?;
Ok(failed == 0) Ok(cnt.get_failure() == 0)
}
pub async fn run_side(
client: Arc<ClientManager>,
queue: SharedQueue<DownloadTask>,
cnt: AtomicCounters,
manifest: Arc<AutoSaveManifest>,
profile: Arc<(String, Arc<RwLock<Profile>>)>,
) -> anyhow::Result<()> {
let p1 = Arc::new(profile.1.read().await.clone());
tokio::fs::create_dir_all(&p1.path).await?;
let sync_resp = sync(&mut client.get_client().await?, &p1).await?;
let id = sync_resp.id;
let n = p1.concurrent.unwrap_or(5);
let semaphore = Arc::new(Semaphore::new(n));
let mut join_set = JoinSet::new();
let cancel_token = CancellationToken::new();
while let Some(task) = queue.try_pop() {
if cancel_token.is_cancelled() {
break;
}
let permit = semaphore.clone().acquire_owned().await?;
if cancel_token.is_cancelled() {
break;
}
let client = client.clone();
let id = id.clone();
let local_manifest = manifest.clone();
let p1 = p1.clone();
let cancel_token = cancel_token.clone();
join_set.spawn(async move {
let guard = task;
let task = &guard.item;
let req = DownloadRequest {
id: id.clone(),
task: task.clone(),
};
let mut conn = client.get_client().await?;
let mut result = download(&mut conn, &req, &p1).await;
if let Err(e) = &result
&& e.downcast_ref::<h2::Error>().is_some()
{
let mut retry_conn = client.get_client().await?;
result = download(&mut retry_conn, &req, &p1).await;
}
if let Err(e) = &result
&& e.downcast_ref::<h2::Error>().is_some()
{
cancel_token.cancel();
}
result?;
local_manifest
.add_bundle(task.bundle_path.clone(), task.bundle_hash.clone())
.await?;
drop(permit);
Ok::<(), anyhow::Error>(())
});
}
while let Some(r) = join_set.join_next().await {
match r {
Ok(Ok(())) => {
cnt.inc_success();
}
Ok(Err(e)) => {
if e.to_string()
.contains("Session did not reconnect within 15s")
{
return Err(anyhow!(e));
}
error!("{}", e);
cnt.inc_failure();
}
Err(e) => {
error!("{}", e);
cnt.inc_failure();
}
}
}
manifest.save().await?;
let req = CloseRequest { id: id.clone() };
close(&mut client.get_client().await?, &req).await?;
Ok(())
} }
#[derive(Debug, Clone, Serialize, Deserialize)] #[derive(Debug, Clone, Serialize, Deserialize)]
@@ -158,3 +283,53 @@ impl AutoSaveManifest {
Ok(()) Ok(())
} }
} }
/// 支持 Clone 的双原子计数器
#[derive(Clone)]
pub struct AtomicCounters {
inner: Arc<CountersInner>,
}
struct CountersInner {
success: AtomicUsize,
failure: AtomicUsize,
}
impl AtomicCounters {
pub fn new() -> Self {
Self {
inner: Arc::new(CountersInner {
success: AtomicUsize::new(0),
failure: AtomicUsize::new(0),
}),
}
}
pub fn inc_success(&self) {
self.inner.success.fetch_add(1, Ordering::Relaxed);
}
pub fn inc_failure(&self) {
self.inner.failure.fetch_add(1, Ordering::Relaxed);
}
pub fn get_success(&self) -> usize {
self.inner.success.load(Ordering::Relaxed)
}
pub fn get_failure(&self) -> usize {
self.inner.failure.load(Ordering::Relaxed)
}
// pub fn load(&self) -> (usize, usize) {
// (
// self.inner.success.load(Ordering::Relaxed),
// self.inner.failure.load(Ordering::Relaxed),
// )
// }
pub fn reset(&self) {
self.inner.success.store(0, Ordering::Relaxed);
self.inner.failure.store(0, Ordering::Relaxed);
}
}
+3 -1
View File
@@ -18,7 +18,8 @@ pub async fn server_send_files<P: AsRef<Path>, S: AsRef<str>>(
let metadata = file.metadata().await?; let metadata = file.metadata().await?;
let file_name = &path_ref.1; let file_name = &path_ref.1;
let name_bytes = file_name.as_ref().as_bytes(); let file_name = file_name.as_ref().replace('\\', "/");
let name_bytes = file_name.as_bytes();
let file_size = metadata.len(); let file_size = metadata.len();
let mut meta_buf = BytesMut::with_capacity(2 + name_bytes.len() + 8); let mut meta_buf = BytesMut::with_capacity(2 + name_bytes.len() + 8);
@@ -62,6 +63,7 @@ pub async fn client_receive(
let mut name_buf = vec![0u8; name_len as usize]; let mut name_buf = vec![0u8; name_len as usize];
response_body_reader.read_exact(&mut name_buf).await?; response_body_reader.read_exact(&mut name_buf).await?;
let file_name = String::from_utf8(name_buf)?; let file_name = String::from_utf8(name_buf)?;
let file_name = file_name.replace('\\', "/");
let data_len = response_body_reader.read_u64().await?; let data_len = response_body_reader.read_u64().await?;
+2
View File
@@ -12,6 +12,8 @@ pub struct SyncContext {
#[serde(default)] #[serde(default)]
pub asset_hash: Option<String>, pub asset_hash: Option<String>,
#[serde(default)] #[serde(default)]
pub app_version: Option<String>,
#[serde(default)]
pub exact_single_file_bundle: bool, pub exact_single_file_bundle: bool,
} }
+4 -1
View File
@@ -17,4 +17,7 @@ webpki-roots = { workspace = true }
anyhow = { workspace = true } anyhow = { workspace = true }
tokio-util = { workspace = true } tokio-util = { workspace = true }
futures-util = { workspace = true } futures-util = { workspace = true }
futures = "0.3.32" futures = { workspace = true }
tokio-socks = { workspace = true }
url = {workspace = true}
async-http-proxy = {workspace = true, features = ["tokio", "runtime-tokio"]}
+2
View File
@@ -34,6 +34,7 @@ pub struct TcpClientTunnelConfig {
pub host: Option<String>, pub host: Option<String>,
pub url: String, pub url: String,
pub token: String, pub token: String,
pub proxy: Option<String>,
} }
impl TcpClientTunnelConfig { impl TcpClientTunnelConfig {
@@ -43,6 +44,7 @@ impl TcpClientTunnelConfig {
host: self.host, host: self.host,
url: self.url, url: self.url,
token: self.token, token: self.token,
proxy: self.proxy,
} }
} }
} }
+91 -12
View File
@@ -1,10 +1,11 @@
use anyhow::anyhow; use anyhow::anyhow;
use async_http_proxy::http_connect_tokio;
use bytes::Bytes; use bytes::Bytes;
use h2::{RecvStream, client, server}; use h2::{RecvStream, client, server};
use http::Request; use http::Request;
use log::{debug, error, info}; use log::{debug, error, info};
use serde::{Deserialize, Serialize}; use serde::{Deserialize, Serialize};
use std::collections::{HashMap, HashSet}; use std::collections::HashMap;
use std::error::Error; use std::error::Error;
use std::fs::File; use std::fs::File;
use std::io::BufReader; use std::io::BufReader;
@@ -17,6 +18,8 @@ use tokio::time::{Duration, sleep, timeout};
use tokio_rustls::rustls; use tokio_rustls::rustls;
use tokio_rustls::rustls::pki_types::{CertificateDer, ServerName}; use tokio_rustls::rustls::pki_types::{CertificateDer, ServerName};
use tokio_rustls::{TlsAcceptor, TlsConnector}; use tokio_rustls::{TlsAcceptor, TlsConnector};
use tokio_socks::tcp::Socks5Stream;
use url::Url;
#[derive(Debug, Copy, Clone, PartialEq, Eq, Serialize, Deserialize)] #[derive(Debug, Copy, Clone, PartialEq, Eq, Serialize, Deserialize)]
pub enum Identity { pub enum Identity {
@@ -39,6 +42,7 @@ pub struct ClientTunnelConfig {
pub url: String, pub url: String,
pub token: String, pub token: String,
pub identity: Identity, pub identity: Identity,
pub proxy: Option<String>,
} }
impl Default for ClientTunnelConfig { impl Default for ClientTunnelConfig {
@@ -48,6 +52,7 @@ impl Default for ClientTunnelConfig {
url: "127.0.0.1:3333".to_string(), url: "127.0.0.1:3333".to_string(),
token: "super_secret_magic_token".to_string(), token: "super_secret_magic_token".to_string(),
identity: Identity::Client, identity: Identity::Client,
proxy: None,
} }
} }
} }
@@ -95,6 +100,30 @@ pub enum TunnelEndpoint {
Server(Arc<ServerManager>), Server(Arc<ServerManager>),
} }
#[derive(Clone)]
pub enum WeakTunnelEndpoint {
Client(std::sync::Weak<ClientManager>),
Server(std::sync::Weak<ServerManager>),
}
impl TunnelEndpoint {
pub fn downgrade(&self) -> WeakTunnelEndpoint {
match self {
Self::Client(c) => WeakTunnelEndpoint::Client(Arc::downgrade(c)),
Self::Server(s) => WeakTunnelEndpoint::Server(Arc::downgrade(s)),
}
}
}
impl WeakTunnelEndpoint {
pub fn upgrade(&self) -> Option<TunnelEndpoint> {
match self {
Self::Client(c) => c.upgrade().map(TunnelEndpoint::Client),
Self::Server(s) => s.upgrade().map(TunnelEndpoint::Server),
}
}
}
pub struct ClientManager { pub struct ClientManager {
pub session_id: AtomicU64, pub session_id: AtomicU64,
pub current_client: Mutex<Option<client::SendRequest<Bytes>>>, pub current_client: Mutex<Option<client::SendRequest<Bytes>>>,
@@ -231,9 +260,9 @@ enum ResumeResult {
pub struct TunnelListener { pub struct TunnelListener {
listener: TcpListener, listener: TcpListener,
config: ServerTunnelConfig, config: ServerTunnelConfig,
pending_plain_sessions: Mutex<HashSet<u64>>, pending_plain_sessions: Mutex<HashMap<u64, std::time::Instant>>,
next_session_id: AtomicU64, next_session_id: AtomicU64,
active_sessions: Mutex<HashMap<u64, TunnelEndpoint>>, active_sessions: Mutex<HashMap<u64, WeakTunnelEndpoint>>,
} }
impl TunnelListener { impl TunnelListener {
@@ -243,7 +272,7 @@ impl TunnelListener {
Ok(Self { Ok(Self {
listener, listener,
config, config,
pending_plain_sessions: Mutex::new(HashSet::new()), pending_plain_sessions: Mutex::new(HashMap::new()),
next_session_id: AtomicU64::new(1), next_session_id: AtomicU64::new(1),
active_sessions: Mutex::new(HashMap::new()), active_sessions: Mutex::new(HashMap::new()),
}) })
@@ -264,7 +293,9 @@ impl TunnelListener {
let mut stream = match self.try_resume_plain_session(stream, peer_addr).await? { let mut stream = match self.try_resume_plain_session(stream, peer_addr).await? {
ResumeResult::NewSession(ep_raw, sid) => { ResumeResult::NewSession(ep_raw, sid) => {
let ep = wrap_raw_endpoint(sid, ep_raw, None); let ep = wrap_raw_endpoint(sid, ep_raw, None);
self.active_sessions.lock().await.insert(sid, ep.clone()); let mut sessions = self.active_sessions.lock().await;
sessions.retain(|_, weak_ep| weak_ep.upgrade().is_some());
sessions.insert(sid, ep.downgrade());
return Ok(ep); return Ok(ep);
} }
ResumeResult::ResumedExisting => continue, ResumeResult::ResumedExisting => continue,
@@ -290,7 +321,9 @@ impl TunnelListener {
let sid = self.next_session_id.fetch_add(1, Ordering::Relaxed); let sid = self.next_session_id.fetch_add(1, Ordering::Relaxed);
let ep = wrap_raw_endpoint(sid, ep_raw, None); let ep = wrap_raw_endpoint(sid, ep_raw, None);
self.active_sessions.lock().await.insert(sid, ep.clone()); let mut sessions = self.active_sessions.lock().await;
sessions.retain(|_, weak_ep| weak_ep.upgrade().is_some());
sessions.insert(sid, ep.downgrade());
return Ok(ep); return Ok(ep);
} }
} }
@@ -320,7 +353,8 @@ impl TunnelListener {
let session_id = self.next_session_id.fetch_add(1, Ordering::Relaxed); let session_id = self.next_session_id.fetch_add(1, Ordering::Relaxed);
{ {
let mut pending = self.pending_plain_sessions.lock().await; let mut pending = self.pending_plain_sessions.lock().await;
pending.insert(session_id); pending.retain(|_, time| time.elapsed() < Duration::from_secs(30));
pending.insert(session_id, std::time::Instant::now());
} }
tls_stream.write_all(TLS_BOOTSTRAP_MAGIC).await?; tls_stream.write_all(TLS_BOOTSTRAP_MAGIC).await?;
@@ -351,7 +385,7 @@ impl TunnelListener {
let is_pending = { let is_pending = {
let mut pending = self.pending_plain_sessions.lock().await; let mut pending = self.pending_plain_sessions.lock().await;
pending.remove(&session_id) pending.remove(&session_id).is_some()
}; };
if is_pending { if is_pending {
@@ -366,7 +400,11 @@ impl TunnelListener {
let ep_raw = upgrade_to_h2_raw(stream, self.config.identity) let ep_raw = upgrade_to_h2_raw(stream, self.config.identity)
.await .await
.map_err(|e| anyhow::anyhow!("{}", e))?; .map_err(|e| anyhow::anyhow!("{}", e))?;
update_endpoint(&ep, ep_raw).await; update_endpoint(
&ep.upgrade().ok_or(anyhow!("Connection is cleared"))?,
ep_raw,
)
.await;
info!( info!(
"[{}] Successfully resumed existing session {}", "[{}] Successfully resumed existing session {}",
peer_addr, session_id peer_addr, session_id
@@ -542,6 +580,47 @@ pub async fn connect_tunnel(config: ClientTunnelConfig) -> Result<TunnelEndpoint
Ok(wrap_raw_endpoint(sid, ep_raw, Some(config))) Ok(wrap_raw_endpoint(sid, ep_raw, Some(config)))
} }
pub async fn connect_with_auto_proxy(config: &ClientTunnelConfig) -> anyhow::Result<TcpStream> {
let Some(proxy) = &config.proxy else {
return Ok(TcpStream::connect(&config.url).await?);
};
let parsed_proxy = Url::parse(proxy)?;
match parsed_proxy.scheme() {
"socks5" => {
let host = parsed_proxy.host_str().unwrap();
let port = parsed_proxy.port().unwrap_or(1080);
let stream = Socks5Stream::connect((host, port), config.url.clone())
.await?
.into_inner();
Ok(stream)
}
"http" | "https" => {
let proxy_addr = format!(
"{}:{}",
parsed_proxy.host_str().unwrap(),
parsed_proxy.port().unwrap_or(80)
);
let mut stream = TcpStream::connect(proxy_addr).await?;
let target_url = Url::parse(&format!("tcp://{}", config.url))?;
http_connect_tokio(
&mut stream,
target_url.host_str().unwrap(),
target_url.port().unwrap_or(80),
)
.await?;
Ok(stream)
}
_ => Err(anyhow::anyhow!(
"Unsupported proxy scheme: {}",
parsed_proxy.scheme()
)),
}
}
async fn do_client_reconnect( async fn do_client_reconnect(
config: &ClientTunnelConfig, config: &ClientTunnelConfig,
current_sid: &mut u64, current_sid: &mut u64,
@@ -563,7 +642,7 @@ async fn do_client_reconnect(
*current_sid = sid; *current_sid = sid;
resume_tunnel_client(config, sid).await resume_tunnel_client(config, sid).await
} else { } else {
let mut stream = TcpStream::connect(&config.url).await?; let mut stream = connect_with_auto_proxy(config).await?;
perform_client_handshake(&mut stream, &config.token, config.identity).await?; perform_client_handshake(&mut stream, &config.token, config.identity).await?;
let raw = upgrade_to_h2_raw(stream, config.identity).await?; let raw = upgrade_to_h2_raw(stream, config.identity).await?;
*current_sid = 0; *current_sid = 0;
@@ -576,7 +655,7 @@ async fn bootstrap_tls_and_get_sid(
host: &str, host: &str,
) -> anyhow::Result<(u64, ())> { ) -> anyhow::Result<(u64, ())> {
let connector = build_client_tls_connector(); let connector = build_client_tls_connector();
let tcp = TcpStream::connect(&config.url).await?; let tcp = connect_with_auto_proxy(config).await?;
let server_name = ServerName::try_from(host.to_string()) let server_name = ServerName::try_from(host.to_string())
.map_err(|_| anyhow!("Invalid TLS host: {}", host))? .map_err(|_| anyhow!("Invalid TLS host: {}", host))?
.to_owned(); .to_owned();
@@ -604,7 +683,7 @@ async fn resume_tunnel_client(
config: &ClientTunnelConfig, config: &ClientTunnelConfig,
session_id: u64, session_id: u64,
) -> anyhow::Result<TunnelEndpointRaw> { ) -> anyhow::Result<TunnelEndpointRaw> {
let mut plain_stream = TcpStream::connect(&config.url).await?; let mut plain_stream = connect_with_auto_proxy(config).await?;
plain_stream.write_all(RESUME_MAGIC).await?; plain_stream.write_all(RESUME_MAGIC).await?;
plain_stream.write_all(&session_id.to_be_bytes()).await?; plain_stream.write_all(&session_id.to_be_bytes()).await?;
upgrade_to_h2_raw(plain_stream, config.identity).await upgrade_to_h2_raw(plain_stream, config.identity).await
+249
View File
@@ -0,0 +1,249 @@
# ==========================================
# Sekai Unpacker 客户端配置示例
# sekai-unpacker-client.example.yaml
# ==========================================
# 日志级别: DEBUG, INFO, WARN, ERROR
log_level: "INFO"
# 说明:`client`/`server` 是 TCP 连接方向,不是进程身份。
# - `client`: 主动连出到远端
# - `server`: 本地监听并接受远端连入
# 当前文件默认展示“client 进程常用连出模式”,也可同时/单独启用 `server`。
# ==========================================
# 服务器连接配置
# ==========================================
# 可配置多个服务器,客户端会依次尝试连接
client:
# 第一个服务器配置
- url: "127.0.0.1:3333" # 服务器地址和端口
token: "your_auth_token_here" # 认证令牌,必须与服务器配置相同
# host: "example.com" # [可选] TLS SNI 主机名
# 仅在需要 TLS 加密到 Identity 交换为止时配置
# 如不配置,则使用明文握手
# 加密范围:Identity 交换之前的握手,H2 通信不加密
# 可以配置多个服务器
# - url: "backup.example.com:3333"
# token: "backup_token"
# host: "backup.example.com"
# ==========================================
# 可选:反向连接场景(本进程作为监听端)
# ==========================================
# server:
# - url: "0.0.0.0:3334" # 本地监听地址和端口
# token: "reverse_link_token" # 需与对端 client.token 一致
# # cert: "/path/to/certificate.pem" # [可选] 启用 TLS 握手
# # key: "/path/to/private.key"
# ==========================================
# 同步配置文件(Profiles
# ==========================================
# 每个配置文件对应一个地区的同步任务
profiles:
# ==========================================
# 国服配置示例
# ==========================================
cn:
# [必需] 地区标识,必须与服务器配置中的 regions 一致
region: cn
# ========== 版本信息配置(可选) ==========
# 以下三个字段会被 dynamic_load 中的值覆盖
# 应用版本
# app_version: "6.0.0"
# 资源版本
# asset_version: "1.0.0"
# 资源哈希值
# asset_hash: "abc123def456"
# ========== 动态加载配置(可选) ==========
# 从指定 URL 动态加载版本和哈希信息
# 用于获取最新版本而无需更新配置文件
dynamic_load:
url: "https://api.example.com/versions/cn"
# 映射关系:本地字段名 -> 远程 JSON 字段名
map:
app_version: appVersion
asset_version: assetVersion
asset_hash: assetHash
# concurrent: concurrent # [可选] 远程并发数,需为数字
# ========== 下载配置 ==========
# 更新间隔(秒,可选)
# 如设置,则每隔指定秒数自动检查更新
# interval: 3600
# 并发下载数(单个资源包的同时下载线程数)
concurrent: 50
# 展开单文件包
exact_single_file_bundle: true
# 导出基础路径
path: "./data/cn"
# ========== 资源过滤配置 ==========
filters:
# 启动应用时必须下载的资源类型
start_app:
- "thumbnail" # 缩略图
- "stamp" # 表情
# - "area" # 地图
# - "home" # 主页背景
# 按需下载(用户请求时才下载)的资源类型
on_demand:
- "thumbnail"
- "stamp"
# - "event" # 活动资源
# - "mysekai" # 我的SEKAI资源
# 跳过下载的资源标识
# 完全不下载这些资源
skip: []
# 按文件扩展名过滤
# 仅下载指定扩展名的文件
file_ext: []
# ========== 导出和处理配置 ==========
export:
# 是否按分类导出(按资源类别分文件夹)
by_category: false
# ---- USM 视频文件处理 ----
usm:
export: true # 是否导出 USM 文件
decode: true # 是否解码为视频格式(需要 FFmpeg)
# ---- ACB 音频包处理 ----
acb:
export: true # 是否导出 ACB 文件
decode: true # 是否解包并解码(需要 cri_ware_decoder
# ---- HCA 音频文件处理 ----
hca:
decode: true # 是否解码为 MP3/FLAC(需要 FFmpeg
# ---- 图片处理 ----
images:
convert_to_webp: false # 是否转换为 WebP 格式(需要 FFmpeg
remove_png: false # 转换后是否删除原 PNG 文件
# ---- 视频处理 ----
video:
convert_to_mp4: false # 是否转换为 MP4 格式(需要 FFmpeg)
direct_usm_to_mp4_with_ffmpeg: false # 直接用 FFmpeg 转换 USM 为 MP4(不先转M2V
remove_m2v: false # 转换后是否删除原 M2V 文件
# ---- 音频处理 ----
audio:
convert_to_mp3: false # 是否转换为 MP3 格式(需要 FFmpeg)
convert_to_flac: false # 是否转换为 FLAC 格式(需要 FFmpeg
remove_wav: false # 转换后是否删除原 WAV 文件
# ==========================================
# 日服配置示例
# ==========================================
jp:
# [必需] 地区标识
region: jp
# 并发下载数
concurrent: 50
# 展开单文件包
exact_single_file_bundle: true
# 导出路径
path: "./data/jp"
# 动态加载配置(使用 GitHub 作为版本源)
dynamic_load:
url: "https://github.com/Team-Haruki/haruki-sekai-master/raw/refs/heads/main/versions/current_version.json"
map:
# app_version: appVersion
asset_version: assetVersion
asset_hash: assetHash
# concurrent: concurrent # [可选] 远程并发数,需为数字
# 过滤配置
filters:
start_app:
- "thumbnail"
- "stamp"
on_demand:
- "thumbnail"
- "stamp"
skip: []
file_ext: []
# 导出配置
export:
by_category: false
usm:
export: true
decode: true
acb:
export: true
decode: true
hca:
decode: true
images:
convert_to_webp: false
remove_png: false
video:
convert_to_mp4: false
direct_usm_to_mp4_with_ffmpeg: false
remove_m2v: false
audio:
convert_to_mp3: false
convert_to_flac: false
remove_wav: false
# ==========================================
# 配置说明
# ==========================================
#
# 【基本配置】
# - log_level: 日志输出级别,建议生产环境设为 INFO
# - client: 主动连出配置,支持故障转移
# - server: 本地监听配置(可选,反向连接时使用)
#
# 【版本管理】
# - 优先级:直接配置(app_version等)> dynamic_load > 默认值
# - 建议使用 dynamic_load 保持版本自动更新
#
# 【资源过滤】
# - start_app: 影响初始化和启动时的下载
# - on_demand: 需要时才下载的资源
# - skip: 完全忽略的资源,不会被下载
# - file_ext: 为空表示下载所有文件,指定则仅下载这些扩展名
#
# 【导出处理】
# - export:true 必须为 true,资源才会被导出
# - decode:true 表示对导出的资源进行解码/转换
# - 各转换功能均需要对应的外部工具(FFmpeg 等)
#
# 【多地区配置】
# - 可在 profiles 中定义多个地区
# - 每个地区独立管理下载和导出配置
# - region 名称必须与服务器配置一致
#
# 【TLS 加密】
# - 仅在需要加密握手时配置 host 字段
# - 加密仅到 Identity 交换,之后 H2 通信不加密
# - 留空或不配置则使用明文握手
#
# 【角色说明】
# - 角色由连接方向决定,不由可执行文件名称决定
# - 同一进程可只开 `client`、只开 `server`,或两者同时开启
+15 -11
View File
@@ -1,19 +1,20 @@
log_level: "DEBUG" log_level: "DEBUG"
#client: client:
# - url: "127.0.0.1:3333" - url: "127.0.0.1:3333"
# token: abc
server:
- url: 127.0.0.1:3333
token: abc token: abc
cert: "D:\\WorkDir\\Nginx\\cert\\_.bluemangoo.net\\_.bluemangoo.net-chain.pem" #server:
key: "D:\\WorkDir\\Nginx\\cert\\_.bluemangoo.net\\_.bluemangoo.net-key.pem" # - url: 127.0.0.1:3333
# token: abc
# cert: "D:\\WorkDir\\Nginx\\cert\\_.bluemangoo.net\\_.bluemangoo.net-chain.pem"
# key: "D:\\WorkDir\\Nginx\\cert\\_.bluemangoo.net\\_.bluemangoo.net-key.pem"
profiles: profiles:
cn: cn:
region: cn region: cn
# interval: 3 # seconds # interval: 3 # seconds
concurrent: 50 concurrent: 50
app_version: "6.0.0"
filters: filters:
start_app: start_app:
- "thumbnail" - "thumbnail"
@@ -44,7 +45,7 @@ profiles:
convert_to_mp3: false convert_to_mp3: false
convert_to_flac: false convert_to_flac: false
remove_wav: false remove_wav: false
# exact_single_file_bundle: true exact_single_file_bundle: true
path: "./data/cn" path: "./data/cn"
jp: jp:
region: jp region: jp
@@ -58,8 +59,11 @@ profiles:
- "stamp" - "stamp"
skip: [ ] skip: [ ]
file_ext: [ ] file_ext: [ ]
asset_version: 6.4.0.30 dynamic_load:
asset_hash: cce60d07-d60e-48dd-be22-12eac2e67950 url: "https://github.com/Team-Haruki/haruki-sekai-master/raw/refs/heads/main/versions/current_version.json"
map:
asset_version: assetVersion
asset_hash: assetHash
export: export:
by_category: false by_category: false
usm: usm:
@@ -81,5 +85,5 @@ profiles:
convert_to_mp3: false convert_to_mp3: false
convert_to_flac: false convert_to_flac: false
remove_wav: false remove_wav: false
# exact_single_file_bundle: true exact_single_file_bundle: true
path: "./data/jp" path: "./data/jp"
+283
View File
@@ -0,0 +1,283 @@
# ==========================================
# Sekai Unpacker 服务器配置示例
# sekai-unpacker-server.example.yaml
# ==========================================
# 日志级别: DEBUG, INFO, WARN, ERROR
log_level: "INFO"
# 说明:`server`/`client` 是 TCP 连接方向,不是进程身份。
# - `server`: 本地监听并接受远端连入
# - `client`: 主动连出到远端
# 当前文件默认展示“server 进程常用监听模式”,也可同时/单独启用 `client`。
# ==========================================
# 服务器监听配置
# ==========================================
# 可配置多个监听地址和端口
server:
# 主服务器配置
- url: "127.0.0.1:3333" # 监听地址和端口
token: "your_secure_token_here" # 认证令牌,对端必须提供相同的令牌
# ========== TLS 证书配置(可选) ==========
# 如配置证书和密钥,服务器将支持 TLS 加密握手
# 如未配置,服务器仅接受明文握手
# 注意:加密仅应用于 Identity 交换阶段,H2 通信保持明文
# cert: "/path/to/certificate.pem" # TLS 证书文件路径
# key: "/path/to/private.key" # TLS 私钥文件路径
# 例子:
# cert: "D:\\certs\\_.example.com-chain.pem"
# key: "D:\\certs\\_.example.com-key.pem"
# 可配置多个监听地址和端口
# - url: "0.0.0.0:3334"
# token: "alternative_token"
# cert: "/path/to/another-cert.pem"
# key: "/path/to/another-key.pem"
# ==========================================
# 可选:反向连接场景(本进程主动连出)
# ==========================================
# client:
# - url: "upstream.example.com:3333" # 远端监听地址和端口
# token: "reverse_link_token" # 需与对端 server.token 一致
# # host: "upstream.example.com" # [可选] 启用 TLS 握手时填写
# ==========================================
# 资源执行配置
# ==========================================
execution:
# HTTP 代理设置(可选)
# 如需通过代理服务器访问资源源,配置此项
# proxy: "http://proxy.example.com:8080"
# proxy: "socks5://proxy.example.com:1080"
# 资源下载和处理的超时时间(秒)
timeout_seconds: 300
# 是否允许客户端取消正在执行的任务
allow_cancel: true
# ========== 重试策略 ==========
retry:
# 最大重试次数(网络错误时)
attempts: 4
# 初始退避时间(毫秒)
# 首次失败后等待 1000ms 再重试
initial_backoff_ms: 1000
# 最大退避时间(毫秒)
# 后续每次重试的等待时间翻倍,但不超过此值
max_backoff_ms: 4000
# ==========================================
# 工具配置
# ==========================================
tools:
# ========== FFmpeg 工具 ==========
# 用于视频和音频的转换和处理
# 留空或 "ffmpeg" 表示使用系统 PATH 中的 ffmpeg
# 如 FFmpeg 未在 PATH 中,请指定完整路径
ffmpeg_path: "ffmpeg"
# 例子:
# ffmpeg_path: "C:\\Program Files\\ffmpeg\\bin\\ffmpeg.exe"
# ffmpeg_path: "/usr/bin/ffmpeg"
# ========== AssetStudio CLI 工具(可选) ==========
# 用于解析 Unity AssetBundle 文件
# 如不配置,某些资源的精细解析功能可能不可用
# asset_studio_cli_path: "D:\\Workspace\\AssetStudio\\AssetStudioCLI\\bin\\Release\\net9.0\\AssetStudioModCLI.exe"
# asset_studio_cli_path: "/opt/AssetStudioCLI"
# ==========================================
# 并发配置
# ==========================================
# 各处理阶段的并发数量配置
# 数值越大,处理速度越快,但对系统资源占用也越大
concurrency:
# 同时下载的资源文件数(好像没用到)
download: 4
# 同时上传/传输的资源文件数(好像没用到)
upload: 4
# ACB 音频包解析的并发数
acb: 8
# USM 视频解析的并发数(好像没用到)
usm: 4
# HCA 音频解码的并发数(通常可设置较大值)
hca: 16
# ==========================================
# 地区配置
# ==========================================
# 定义支持的资源地区及其对应的数据源
regions:
# ==========================================
# 日本服务器配置
# ==========================================
jp:
# 是否启用该地区
enabled: true
# ========== 资源提供者配置 ==========
provider:
# 提供者类型:colorful_palette 或 nuverse
kind: colorful_palette
# 资源包元信息 API 地址模板
# {env}:环境标识(production/staging 等)
# {hash}:配置哈希值
# {asset_version}:资源版本
# {asset_hash}:资源包哈希
asset_info_url_template: "https://{env}-{hash}-assetbundle-info.sekai.colorfulpalette.org/api/version/{asset_version}/{asset_hash}/os/ios"
# 资源包下载地址模板
# {bundle_path}:资源包路径
asset_bundle_url_template: "https://{env}-{hash}-assetbundle.sekai.colorfulpalette.org/{asset_version}/{asset_hash}/ios/{bundle_path}"
# 配置文件名(如 production
profile: "production"
# 配置对应的哈希值映射表
profile_hashes:
production: "cf2d2388"
# staging: "12345678"
# 是否需要 Cookie
# 某些源需要 Cookie 才能访问(防盗链)
required_cookies: true
# 获取 Cookie 的 Bootstrap URL(可选)
# 如配置,服务器会先访问此 URL 获取 Cookie
# cookie_bootstrap_url: "https://example.com/bootstrap"
# ========== 加密配置 ==========
# 资源包的 AES 加密参数
crypto:
# AES 密钥(十六进制字符串)
aes_key_hex: "6732666343305a637a4e394d544a3631"
# AES 初始向量(十六进制字符串)
aes_iv_hex: "6d737833495630693958453575595a31"
# ========== 运行时配置 ==========
runtime:
# 资源对应的 Unity 版本
# 用于正确解析 AssetBundle
unity_version: "2022.3.21f1"
# ==========================================
# 中国服务器配置
# ==========================================
cn:
enabled: true
provider:
# Nuverse 提供者(国服使用)
kind: nuverse
# 获取资源版本的 API 地址
# {app_version}:应用版本
asset_version_url: "https://lf3-mkcncdn-tos.dailygn.com/obj/rt-game-lf/gdl_app_5236/Mainland/{app_version}/Release/cn_online/ios/version"
# 资源包元信息 API 地址模板
# {asset_version}:资源版本
asset_info_url_template: "https://lf3-mkcncdn-tos.dailygn.com/obj/sf-game-lf/gdl_app_5236/AssetBundle/{app_version}/Release/cn_online/ios{asset_version}/AssetBundleInfoNew.json"
# 资源包下载地址模板
# {app_version}:应用版本
# {bundle_path}:资源包路径
asset_bundle_url_template: "https://lf3-mkcncdn-tos.dailygn.com/obj/sf-game-lf/gdl_app_5236/AssetBundle/{app_version}/Release/cn_online/{bundle_path}"
# 是否需要 Cookie
required_cookies: false
# cookie_bootstrap_url: null
crypto:
aes_key_hex: "6732666343305a637a4e394d544a3631"
aes_iv_hex: "6d737833495630693958453575595a31"
runtime:
unity_version: "2022.3.21f1"
# ==========================================
# 其他地区配置示例(禁用)
# ==========================================
# tw:
# enabled: false
# provider:
# kind: colorful_palette
# asset_info_url_template: "https://{env}-{hash}-assetbundle-info.sekai.colorfulpalette.org/api/version/{asset_version}/{asset_hash}/os/ios"
# asset_bundle_url_template: "https://{env}-{hash}-assetbundle.sekai.colorfulpalette.org/{asset_version}/{asset_hash}/ios/{bundle_path}"
# profile: "production"
# profile_hashes:
# production: "cf2d2388"
# required_cookies: true
# crypto:
# aes_key_hex: "6732666343305a637a4e394d544a3631"
# aes_iv_hex: "6d737833495630693958453575595a31"
# runtime:
# unity_version: "2022.3.21f1"
# ==========================================
# 配置说明
# ==========================================
#
# 【服务器配置】
# - 可配置多个监听地址,支持多个客户端连接
# - 每个监听配置需要不同的 token 来区分认证
# - cert 和 key 配置是可选的:
# * 配置了:支持 TLS 加密握手 + 明文握手(双模式)
# * 未配置:仅支持明文握手
# - 如需主动连出,可额外配置 `client` 区块
#
# 【认证机制】
# - client 中的 token 必须与某个 server 配置中的 token 匹配
# - 握手时会验证 token,不匹配则拒绝连接
#
# 【TLS 加密】
# - 加密范围:Identity 交换(身份验证)阶段
# - H2 通信(实际数据传输)始终不加密,基于 TCP 流
# - 证书应该是 X.509 格式的 PEM 编码文件
#
# 【资源源配置】
# - 支持两种提供者:
# * colorful_palette:日服资源源
# * nuverse:国服资源源
# - 每个提供者有不同的 API 地址格式
#
# 【并发控制】
# - 合理设置并发数可提升处理速度
# - 过高可能导致网络压力或系统资源不足
# - 建议:download/upload 4-8hca 16,其他 4-8
#
# 【版本信息】
# - 客户端会根据 region 字段从服务器获取该地区的配置
# - 确保客户端 profiles 中的 region 值与服务器的地区名称一致
#
# 【禁用地区】
# - 将 enabled 设为 false 可禁用某个地区
# - 禁用的地区客户端无法同步
#
# 【代理配置】
# - 如果服务器需要通过代理访问资源源,在 execution.proxy 中配置
# - 支持 HTTP、HTTPS、SOCKS5 代理
#
# 【工具路径】
# - 确保 ffmpeg 和 AssetStudioCLI 可从配置的路径访问
# - 如在 PATH 中,可直接写命令名称
# - 否则需要提供完整的绝对路径
#
# 【角色说明】
# - 角色由连接方向决定,不由可执行文件名称决定
# - 同一进程可只开 `client`、只开 `server`,或两者同时开启
+6 -7
View File
@@ -1,12 +1,12 @@
log_level: "DEBUG" log_level: "DEBUG"
client: #client:
- url: "127.0.0.1:3333" # - url: "127.0.0.1:3333"
token: abc
host: "local.bluemangoo.net"
#server:
# - url: 127.0.0.1:3333
# token: abc # token: abc
# host: "local.bluemangoo.net"
server:
- url: 127.0.0.1:3333
token: abc
execution: execution:
proxy: "" proxy: ""
@@ -49,7 +49,6 @@ regions:
provider: provider:
kind: nuverse kind: nuverse
asset_version_url: "https://lf3-mkcncdn-tos.dailygn.com/obj/rt-game-lf/gdl_app_5236/Mainland/{app_version}/Release/cn_online/ios/version" asset_version_url: "https://lf3-mkcncdn-tos.dailygn.com/obj/rt-game-lf/gdl_app_5236/Mainland/{app_version}/Release/cn_online/ios/version"
app_version: "5.2.0"
asset_info_url_template: "https://lf3-mkcncdn-tos.dailygn.com/obj/sf-game-lf/gdl_app_5236/AssetBundle/{app_version}/Release/cn_online/ios{asset_version}/AssetBundleInfoNew.json" asset_info_url_template: "https://lf3-mkcncdn-tos.dailygn.com/obj/sf-game-lf/gdl_app_5236/AssetBundle/{app_version}/Release/cn_online/ios{asset_version}/AssetBundleInfoNew.json"
asset_bundle_url_template: "https://lf3-mkcncdn-tos.dailygn.com/obj/sf-game-lf/gdl_app_5236/AssetBundle/{app_version}/Release/cn_online/{bundle_path}" asset_bundle_url_template: "https://lf3-mkcncdn-tos.dailygn.com/obj/sf-game-lf/gdl_app_5236/AssetBundle/{app_version}/Release/cn_online/{bundle_path}"
required_cookies: false required_cookies: false
+43
View File
@@ -0,0 +1,43 @@
FROM rust:slim-bookworm AS builder
WORKDIR /usr/src/app
RUN cargo build --release
RUN apt-get update && apt-get install -y musl-tools pkg-config libssl-dev
RUN rustup target add x86_64-unknown-linux-musl
COPY . .
RUN cargo build --target x86_64-unknown-linux-musl --release --bin client
FROM mcr.microsoft.com/dotnet/sdk:9.0-bookworm-slim AS assetstudio-builder
WORKDIR /src
RUN apt-get update && apt-get install -y --no-install-recommends git ca-certificates && \
rm -rf /var/lib/apt/lists/*
RUN git clone --depth 1 --single-branch --branch sekai-modify https://github.com/Team-Haruki/AssetStudio.git
RUN cd AssetStudio/AssetStudioCLI && \
dotnet publish -c Release -r linux-x64 -f net9.0 --self-contained true -o /app/assetstudio \
-p:PublishTrimmed=false \
-p:PublishSingleFile=true \
-p:IncludeNativeLibrariesForSelfExtract=true
FROM mwader/static-ffmpeg:7.1.1 AS ffmpeg-builder
FROM debian:trixie-slim
RUN apt-get update && apt-get install -y --no-install-recommends \
ca-certificates \
tzdata \
libicu76 \
libxml2 && \
rm -rf /var/lib/apt/lists/*
WORKDIR /app
COPY --from=builder /app/target/release/server /app/server
COPY --from=assetstudio-builder /app/assetstudio /app/server
COPY --from=ffmpeg-builder /ffmpeg /usr/local/bin/ffmpeg
RUN ln -sf /app/assetstudio/AssetStudioModCLI /app/assetstudio/AssetStudioCLI && \
mkdir -p logs
ENV TZ=Asia/Shanghai \
DOTNET_SYSTEM_GLOBALIZATION_INVARIANT=false \
EXPOSE 3000
CMD ["./server"]
+31 -51
View File
@@ -7,6 +7,7 @@ use communicator::http::{json_from_request, send, send_error};
use h2::RecvStream; use h2::RecvStream;
use h2::server::SendResponse; use h2::server::SendResponse;
use http::{Request, Response}; use http::{Request, Response};
use log::debug;
pub async fn download( pub async fn download(
mut request: Request<RecvStream>, mut request: Request<RecvStream>,
@@ -43,58 +44,37 @@ pub async fn download(
} }
let (dir, single_file) = dir.unwrap(); let (dir, single_file) = dir.unwrap();
let dir_base = std::env::temp_dir() let files = if context.sync_context.filters.file_ext.is_empty() {
.join("sekai-updater") find_files(&dir)
.join("extract")
.join(&context.sync_context.region);
let files = if !dir.is_dir() {
if context.sync_context.filters.file_ext.is_empty()
|| dir.extension().is_some_and(|t| {
context
.sync_context
.filters
.file_ext
.contains(&t.to_str().unwrap().to_lowercase())
})
{
vec![(
dir.clone(),
dir.strip_prefix(&dir_base)
.unwrap()
.to_string_lossy()
.to_string(),
)]
} else {
vec![]
}
} else { } else {
let files = if context.sync_context.filters.file_ext.is_empty() { find_files_by_extensions(&dir, &context.sync_context.filters.file_ext)
find_files(&dir)
} else {
find_files_by_extensions(&dir, &context.sync_context.filters.file_ext)
};
if let Err(error) = files {
send_error(send_response, error.into());
return Ok(());
}
let base = dir.strip_prefix(&dir_base).unwrap();
let base_str = if single_file {
base.parent().unwrap().to_str().unwrap()
} else {
base.to_str().unwrap()
};
files
.unwrap()
.iter()
.map(|p| {
let raw = p.clone();
let name = p.strip_prefix(&dir).unwrap();
let path = format!("{}/{}", base_str, name.to_string_lossy());
(raw, path)
})
.collect::<Vec<_>>()
}; };
if let Err(error) = files {
send_error(send_response, error.into());
return Ok(());
}
let files = files
.unwrap()
.iter()
.map(|p| {
let raw = p.clone();
let name = p.strip_prefix(&dir).unwrap();
debug!("{} {}", p.to_string_lossy(), single_file);
let path = if context.sync_context.exact_single_file_bundle
&& single_file
{
format!(
"{}/{}",
name.parent().unwrap().parent().unwrap().to_string_lossy(),
p.file_name().unwrap().to_string_lossy()
)
} else {
name.to_string_lossy().to_string()
};
(raw, path)
})
.collect::<Vec<_>>();
let response = Response::builder() let response = Response::builder()
.status(200) .status(200)
@@ -106,7 +86,7 @@ pub async fn download(
let _ = server_send_files(send_stream, &files).await; let _ = server_send_files(send_stream, &files).await;
} }
let _ = std::fs::remove_file(dir); let _ = std::fs::remove_dir_all(dir);
Ok(()) Ok(())
} }