Files
BlueArchiveToolkit/infrastructure/src/bin/bat_official_sync.rs
T
2026-07-13 22:02:11 +08:00

4134 lines
143 KiB
Rust
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
use bat_adapters::official::yostar_jp::PatchPlatform;
use bat_infrastructure::{
OfficialServerInfoSource, OfficialUpdateConfig, OfficialUpdateProgress, OfficialUpdateReport,
OfficialUpdateService, OfficialUpdateStatus, OfficialVerificationSummary,
};
use serde::{Deserialize, Serialize};
use std::env;
use std::fs::{self, OpenOptions};
use std::io::{BufRead, BufReader, Write};
#[cfg(unix)]
use std::os::unix::net::{UnixListener, UnixStream};
use std::path::{Path, PathBuf};
use std::process::{Command, Stdio};
use std::sync::{Arc, Condvar, Mutex};
use std::thread;
use std::time::{Duration, Instant, SystemTime, UNIX_EPOCH};
const EXIT_ERROR: i32 = 1;
const EXIT_LOCKED: i32 = 75;
const DEFAULT_WATCH_INTERVAL_SECONDS: u64 = 60 * 60;
const DEFAULT_ERROR_RETRY_SECONDS: u64 = 60;
const DEFAULT_DAEMON_STATE_DIR: &str = "/tmp/bat-pid";
const DAEMON_PID_FILE: &str = "bat.pid";
const DAEMON_STATUS_FILE: &str = "bat-status.json";
const DAEMON_LOG_FILE: &str = "bat-daemon.log";
const DAEMON_SOCKET_FILE: &str = "bat.sock";
const DAEMON_CONTROL_LOCK_FILE: &str = "bat-control.lock";
const DAEMON_STATUS_VERSION: u32 = 1;
const SECONDS_PER_DAY: u64 = 24 * 60 * 60;
const BEIJING_UTC_OFFSET_SECONDS: u64 = 8 * 60 * 60;
const DAILY_FORCED_REFRESH_LOCAL_SECONDS: [u64; 3] = [3 * 60 * 60, 16 * 60 * 60, 18 * 60 * 60];
const DAILY_FORCED_REFRESH_LABEL: &str = "UTC+8 03:00, 16:00, 18:00";
const STARTUP_BANNER: &str = r#"
============================================================
____ _ _ _ _ _____ _ _ _ _
| __ )| |_ _ ___ / \ _ __ ___| |__ (_)_ _____|_ _|__ ___ | | | _(_) |_
| _ \| | | | |/ _ \/ _ \ | '__/ __| '_ \| \ \ / / _ \ | |/ _ \ / _ \| | |/ / | __|
| |_) | | |_| | __/ ___ \| | | (__| | | | |\ V / __/ | | (_) | (_) | | <| | |_
|____/|_|\__,_|\__/_/ \_\_| \___|_| |_|_|\_/ \___/ |_|\___/ \___/|_|_|\_\_|\__|
BlueArchiveToolkit
Official Resource Sync
============================================================
"#;
fn main() {
match run() {
Ok(exit_code) if exit_code != 0 => std::process::exit(exit_code),
Ok(_) => {}
Err(error) => {
let error_message = error.to_string();
let exit_code = if is_locked_error(&error_message) {
EXIT_LOCKED
} else {
EXIT_ERROR
};
let payload = ErrorReport {
status: "error",
exit_code,
error: localize_error_message(&error_message),
next_retry_seconds: None,
};
eprintln!(
"{}",
serde_json::to_string_pretty(&payload)
.unwrap_or_else(|_| "{\"status\":\"error\"}".to_string())
);
std::process::exit(exit_code);
}
}
}
fn run() -> anyhow::Result<i32> {
let options = parse_args()?;
if matches!(
options.command,
CliCommand::Run | CliCommand::Refresh | CliCommand::Verify | CliCommand::Repair
) && !options.daemon
&& options.banner
{
print_startup_banner();
}
match options.command {
CliCommand::Run => {
if options.daemon {
let _control_lock = DaemonControlLock::acquire(&options.state_dir)?;
run_daemon_start(options)?;
} else if options.watch {
if !options.daemon_child {
assert_no_live_daemon_output_conflict(&options, "watch")?;
}
run_watch(options)?;
} else {
assert_no_live_daemon_output_conflict(&options, "run")?;
let mut logger = ProgressLogger::new(options.progress);
let report =
OfficialUpdateService::new().run_with_progress(&options.config, |event| {
logger.log(event);
})?;
if should_print_status(report.update_status, options.quiet_up_to_date) {
print_report(options.output_format, &report)?;
}
}
Ok(0)
}
CliCommand::Status => {
let _control_lock = DaemonControlLock::acquire(&options.state_dir)?;
print_daemon_status(&options.state_dir, options.output_format)?;
Ok(0)
}
CliCommand::Stop => {
let _control_lock = DaemonControlLock::acquire(&options.state_dir)?;
stop_daemon(&options.state_dir, options.output_format)?;
Ok(0)
}
CliCommand::Restart => {
let _control_lock = DaemonControlLock::acquire(&options.state_dir)?;
run_daemon_restart(&options, "restart")?;
Ok(0)
}
CliCommand::Reload => {
let _control_lock = DaemonControlLock::acquire(&options.state_dir)?;
run_daemon_restart(&options, "reload")?;
Ok(0)
}
CliCommand::Refresh => {
run_sync_command(&options, "refresh")?;
Ok(0)
}
CliCommand::Verify => {
let healthy = run_verify_command(&options)?;
Ok(if healthy { 0 } else { EXIT_ERROR })
}
CliCommand::Repair => {
run_sync_command(&options, "repair")?;
Ok(0)
}
CliCommand::Logs => {
let _control_lock = DaemonControlLock::acquire(&options.state_dir)?;
run_logs_command(&options)?;
Ok(0)
}
CliCommand::Doctor => {
let healthy = run_doctor_command(&options)?;
Ok(if healthy { 0 } else { EXIT_ERROR })
}
CliCommand::CleanStable => {
run_clean_stable_command(&options)?;
Ok(0)
}
}
}
#[derive(Debug, Serialize)]
struct ErrorReport<'a> {
status: &'a str,
exit_code: i32,
error: String,
#[serde(skip_serializing_if = "Option::is_none")]
next_retry_seconds: Option<u64>,
}
#[derive(Debug, Clone, PartialEq, Eq)]
struct CliOptions {
config: OfficialUpdateConfig,
command: CliCommand,
output_format: OutputFormat,
watch: bool,
daemon: bool,
daemon_child: bool,
state_dir: PathBuf,
output_explicit: bool,
sync_option_explicit: bool,
interval: Duration,
error_retry_interval: Duration,
quiet_up_to_date: bool,
quiet_up_to_date_explicit: bool,
progress: bool,
banner: bool,
tail_lines: usize,
}
impl Default for CliOptions {
fn default() -> Self {
Self {
config: OfficialUpdateConfig::default(),
command: CliCommand::Run,
output_format: OutputFormat::Human,
watch: false,
daemon: false,
daemon_child: false,
state_dir: PathBuf::from(DEFAULT_DAEMON_STATE_DIR),
output_explicit: false,
sync_option_explicit: false,
interval: Duration::from_secs(DEFAULT_WATCH_INTERVAL_SECONDS),
error_retry_interval: Duration::from_secs(DEFAULT_ERROR_RETRY_SECONDS),
quiet_up_to_date: false,
quiet_up_to_date_explicit: false,
progress: true,
banner: true,
tail_lines: 200,
}
}
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum OutputFormat {
Human,
Json,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum CliCommand {
Run,
Status,
Stop,
Restart,
Reload,
Refresh,
Verify,
Repair,
Doctor,
Logs,
CleanStable,
}
fn run_watch(options: CliOptions) -> anyhow::Result<()> {
if options.config.dry_run {
return Err(anyhow::anyhow!("watch 常驻模式不能和 --dry-run 同时使用"));
}
if options.interval.is_zero() {
return Err(anyhow::anyhow!("watch 检查间隔必须大于 0"));
}
if options.error_retry_interval.is_zero() {
return Err(anyhow::anyhow!("watch 失败重试间隔必须大于 0"));
}
let service = OfficialUpdateService::new();
let mut logger = ProgressLogger::new(options.progress);
let daemon_state_dir = options.state_dir.clone();
let daemon_control = if options.daemon_child {
Some(new_daemon_control())
} else {
None
};
let _rpc_server = if let Some(control) = daemon_control.as_ref() {
Some(start_daemon_rpc_server(
&daemon_state_dir,
Arc::clone(control),
)?)
} else {
None
};
let mut next_forced_refresh_at = next_forced_refresh_at_or_after(SystemTime::now());
let mut pending_scheduled_force = false;
let mut pending_rpc_force = false;
logger.log_text(
"watch",
format!(
"每日强制刷新时间:{DAILY_FORCED_REFRESH_LABEL};距离下一次强制刷新还有 {}",
format_duration(duration_until(next_forced_refresh_at, SystemTime::now()))
),
);
loop {
if daemon_control_stop_requested(daemon_control.as_ref()) {
logger.log_text("daemon", "收到停止请求,watch 循环准备退出");
break;
}
let now = SystemTime::now();
if now >= next_forced_refresh_at {
pending_scheduled_force = true;
logger.log_text(
"watch",
format!("已触发固定时间强制刷新:{DAILY_FORCED_REFRESH_LABEL}"),
);
next_forced_refresh_at = next_forced_refresh_after(now);
}
let mut iteration_config = options.config.clone();
if pending_scheduled_force {
iteration_config.force = true;
}
let rpc_force_this_round = pending_rpc_force;
if rpc_force_this_round {
iteration_config.force = true;
pending_rpc_force = false;
}
let mut sleep_for = options.interval;
logger.log_text(
"watch",
if pending_scheduled_force && rpc_force_this_round {
"开始执行 watch 轮次:本轮为固定时间和 RPC 触发的强制刷新"
} else if pending_scheduled_force {
"开始执行 watch 轮次:本轮为固定时间强制刷新"
} else if rpc_force_this_round {
"开始执行 watch 轮次:本轮为 RPC 触发的强制刷新"
} else {
"开始执行 watch 轮次"
},
);
record_daemon_status(
options.daemon_child,
&daemon_state_dir,
&mut logger,
DaemonStatusUpdate {
state: "running",
last_update_status: None,
last_error: None,
next_retry_seconds: None,
pending_scheduled_force,
next_forced_refresh_at,
},
);
match service.run_with_progress_and_cancellation(
&iteration_config,
|event| logger.log(event),
|| daemon_control_stop_requested(daemon_control.as_ref()),
) {
Ok(report) => {
if pending_scheduled_force {
pending_scheduled_force = false;
}
if should_print_status(report.update_status, options.quiet_up_to_date) {
print_report(options.output_format, &report)?;
}
sleep_for =
sleep_for.min(duration_until(next_forced_refresh_at, SystemTime::now()));
record_daemon_status(
options.daemon_child,
&daemon_state_dir,
&mut logger,
DaemonStatusUpdate {
state: "sleeping",
last_update_status: Some(report.update_status.as_str().to_string()),
last_error: None,
next_retry_seconds: Some(sleep_for.as_secs()),
pending_scheduled_force,
next_forced_refresh_at,
},
);
logger.log_text(
"watch",
format!(
"本轮完成:状态={};下次检查将在 {} 后执行;距离下一次固定强制刷新还有 {}",
report.update_status.as_str(),
format_duration(sleep_for),
format_duration(duration_until(next_forced_refresh_at, SystemTime::now()))
),
);
}
Err(error) => {
sleep_for = options.error_retry_interval;
sleep_for =
sleep_for.min(duration_until(next_forced_refresh_at, SystemTime::now()));
record_daemon_status(
options.daemon_child,
&daemon_state_dir,
&mut logger,
DaemonStatusUpdate {
state: "error",
last_update_status: None,
last_error: Some(localize_error_message(&error.to_string())),
next_retry_seconds: Some(sleep_for.as_secs()),
pending_scheduled_force,
next_forced_refresh_at,
},
);
logger.log_text(
"watch",
format!(
"本轮失败;将在 {} 后重试:{}",
format_duration(sleep_for),
localize_error_message(&error.to_string())
),
);
let payload = ErrorReport {
status: "error",
exit_code: EXIT_ERROR,
error: localize_error_message(&error.to_string()),
next_retry_seconds: Some(sleep_for.as_secs()),
};
eprintln!("{}", serde_json::to_string_pretty(&payload)?);
if daemon_control_stop_requested(daemon_control.as_ref()) {
logger.log_text("daemon", "停止请求已中断本轮同步,watch 循环准备退出");
break;
}
}
}
match wait_for_daemon_wake(daemon_control.as_ref(), sleep_for) {
DaemonWake::Timeout => {}
DaemonWake::Stop => {
logger.log_text("daemon", "收到停止请求,watch 循环准备退出");
break;
}
DaemonWake::Refresh { force } => {
pending_rpc_force = force;
logger.log_text(
"daemon",
if force {
"收到 RPC refresh --force 请求,将尽快执行强制刷新"
} else {
"收到 RPC refresh 请求,将尽快执行刷新检查"
},
);
}
DaemonWake::Reload => {
pending_rpc_force = true;
logger.log_text(
"daemon",
"收到 RPC reload 请求,将尽快重新执行自动发现和强制刷新",
);
}
}
}
record_daemon_status(
options.daemon_child,
&daemon_state_dir,
&mut logger,
DaemonStatusUpdate {
state: "stopped",
last_update_status: None,
last_error: None,
next_retry_seconds: None,
pending_scheduled_force,
next_forced_refresh_at,
},
);
if options.daemon_child {
let _ = fs::remove_file(daemon_pid_path(&daemon_state_dir));
let _ = fs::remove_file(daemon_socket_path(&daemon_state_dir));
}
Ok(())
}
#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
struct DaemonStatusFile {
version: u32,
pid: u32,
state: String,
resource_output_root: PathBuf,
state_dir: PathBuf,
log_path: PathBuf,
started_unix_seconds: u64,
updated_unix_seconds: u64,
last_update_status: Option<String>,
last_error: Option<String>,
next_retry_seconds: Option<u64>,
pending_scheduled_force: bool,
next_forced_refresh_unix_seconds: Option<u64>,
command: Vec<String>,
}
#[derive(Debug)]
struct DaemonStatusUpdate<'a> {
state: &'a str,
last_update_status: Option<String>,
last_error: Option<String>,
next_retry_seconds: Option<u64>,
pending_scheduled_force: bool,
next_forced_refresh_at: SystemTime,
}
type DaemonControl = Arc<(Mutex<DaemonControlState>, Condvar)>;
#[derive(Debug, Default)]
struct DaemonControlState {
stop_requested: bool,
refresh_requested: bool,
force_refresh_requested: bool,
reload_requested: bool,
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum DaemonWake {
Timeout,
Stop,
Refresh { force: bool },
Reload,
}
#[derive(Debug, Deserialize)]
struct JsonRpcRequest {
#[allow(dead_code)]
jsonrpc: Option<String>,
id: Option<serde_json::Value>,
method: String,
params: Option<serde_json::Value>,
}
#[derive(Debug, Serialize)]
struct JsonRpcResponse {
jsonrpc: &'static str,
#[serde(skip_serializing_if = "Option::is_none")]
id: Option<serde_json::Value>,
#[serde(skip_serializing_if = "Option::is_none")]
result: Option<serde_json::Value>,
#[serde(skip_serializing_if = "Option::is_none")]
error: Option<JsonRpcError>,
}
#[derive(Debug, Serialize, Deserialize)]
struct JsonRpcError {
code: i32,
message: String,
}
#[derive(Debug, Deserialize)]
struct JsonRpcClientResponse {
#[allow(dead_code)]
jsonrpc: String,
#[allow(dead_code)]
id: Option<serde_json::Value>,
result: Option<serde_json::Value>,
error: Option<JsonRpcError>,
}
#[derive(Debug, Serialize)]
struct DaemonRpcAck {
command: &'static str,
status: &'static str,
message: &'static str,
state_dir: PathBuf,
socket_path: PathBuf,
#[serde(skip_serializing_if = "Option::is_none")]
force: Option<bool>,
}
const RPC_METHOD_STATUS: &str = "bat.status";
const RPC_METHOD_STOP: &str = "bat.stop";
const RPC_METHOD_RELOAD: &str = "bat.reload";
const RPC_METHOD_REFRESH: &str = "bat.refresh";
const RPC_METHOD_LOGS: &str = "bat.logs";
fn record_daemon_status(
enabled: bool,
state_dir: &Path,
logger: &mut ProgressLogger,
update: DaemonStatusUpdate<'_>,
) {
if !enabled {
return;
}
if let Err(error) = update_daemon_status(
state_dir,
update.state,
update.last_update_status,
update.last_error,
update.next_retry_seconds,
update.pending_scheduled_force,
update.next_forced_refresh_at,
) {
logger.log_text("daemon", format!("更新后台状态失败:{error}"));
}
}
#[derive(Debug)]
struct DaemonRpcServer {
socket_path: PathBuf,
}
impl Drop for DaemonRpcServer {
fn drop(&mut self) {
let _ = fs::remove_file(&self.socket_path);
}
}
#[derive(Debug)]
struct DaemonControlLock {
path: PathBuf,
pid: u32,
}
impl DaemonControlLock {
fn acquire(state_dir: &Path) -> anyhow::Result<Self> {
fs::create_dir_all(state_dir)?;
let path = daemon_control_lock_path(state_dir);
let pid = std::process::id();
for attempt in 0..=1 {
match OpenOptions::new().write(true).create_new(true).open(&path) {
Ok(mut file) => {
file.write_all(pid.to_string().as_bytes())?;
return Ok(Self { path, pid });
}
Err(error) if error.kind() == std::io::ErrorKind::AlreadyExists => {
if attempt == 0 && remove_recoverable_pid_lock(&path)? {
continue;
}
return Err(anyhow::anyhow!(
"后台控制命令已被锁定 (locked):{};{};如确认没有 bat 控制命令正在运行,可执行 clean-stable 清理",
path.display(),
describe_pid_lock_owner(&path)?
));
}
Err(error) => {
return Err(anyhow::anyhow!(
"获取后台控制锁失败 {}{error}",
path.display()
));
}
}
}
Err(anyhow::anyhow!("获取后台控制锁失败 {}", path.display()))
}
}
impl Drop for DaemonControlLock {
fn drop(&mut self) {
let expected = self.pid.to_string();
if fs::read_to_string(&self.path)
.map(|contents| contents.trim() == expected)
.unwrap_or(false)
{
let _ = fs::remove_file(&self.path);
}
}
}
#[derive(Debug, Clone, PartialEq, Eq)]
enum PidLockState {
Missing,
Active(u32),
StalePid(u32),
Corrupt,
}
fn classify_pid_lock_file(path: &Path) -> anyhow::Result<PidLockState> {
let contents = match fs::read_to_string(path) {
Ok(contents) => contents,
Err(error) if error.kind() == std::io::ErrorKind::NotFound => {
return Ok(PidLockState::Missing);
}
Err(error) => return Err(error.into()),
};
let Some(pid) = parse_pid_value(&contents) else {
return Ok(PidLockState::Corrupt);
};
if process_exists(pid) {
Ok(PidLockState::Active(pid))
} else {
Ok(PidLockState::StalePid(pid))
}
}
fn remove_recoverable_pid_lock(path: &Path) -> anyhow::Result<bool> {
match classify_pid_lock_file(path)? {
PidLockState::Missing => Ok(false),
PidLockState::Active(_) => Ok(false),
PidLockState::StalePid(_) | PidLockState::Corrupt => {
fs::remove_file(path)?;
Ok(true)
}
}
}
fn describe_pid_lock_owner(path: &Path) -> anyhow::Result<String> {
Ok(match classify_pid_lock_file(path)? {
PidLockState::Missing => "锁文件已不存在".to_string(),
PidLockState::Active(pid) => format!("owner_pid={pid} 仍在运行"),
PidLockState::StalePid(pid) => format!("owner_pid={pid} 已失效"),
PidLockState::Corrupt => "锁文件内容不是有效 PID".to_string(),
})
}
fn parse_pid_value(contents: &str) -> Option<u32> {
contents.trim().parse::<u32>().ok().filter(|pid| *pid > 0)
}
#[derive(Debug, Clone, PartialEq, Eq)]
struct DaemonResourceConflict {
pid: Option<u32>,
daemon_root: Option<PathBuf>,
target_root: PathBuf,
}
fn assert_no_live_daemon_output_conflict(
options: &CliOptions,
command_name: &str,
) -> anyhow::Result<()> {
if options.daemon || options.daemon_child || options.config.dry_run {
return Ok(());
}
let Some(conflict) =
live_daemon_resource_conflict(&options.state_dir, &options.config.output_root)?
else {
return Ok(());
};
let daemon_root = conflict
.daemon_root
.as_ref()
.map(|path| path.display().to_string())
.unwrap_or_else(|| "未知".to_string());
let pid = conflict
.pid
.map(|pid| pid.to_string())
.unwrap_or_else(|| "未知".to_string());
Err(anyhow::anyhow!(
"后台同步进程正在管理同一资源目录 (locked)command={command_name} pid={pid} daemon_output={} target_output={};请通过后台 refresh/reload 控制刷新,或先执行 stop 再手动写入资源",
daemon_root,
conflict.target_root.display()
))
}
fn live_daemon_resource_conflict(
state_dir: &Path,
output_root: &Path,
) -> anyhow::Result<Option<DaemonResourceConflict>> {
let pid_path = daemon_pid_path(state_dir);
let status_path = daemon_status_path(state_dir);
let pid_from_file = read_pid_file(&pid_path)?;
let rpc_available = daemon_rpc_available(state_dir);
let pid_file_live = pid_from_file.map(process_exists).unwrap_or(false);
let status_file = match read_daemon_status_file(&status_path) {
Ok(status_file) => status_file,
Err(error) if pid_file_live || rpc_available => {
return Err(anyhow::anyhow!(
"后台同步进程可能仍在运行,但状态文件不可解析 (locked){}{error};请先执行 status/stop 或 clean-stable 恢复状态目录",
status_path.display()
));
}
Err(_) => return Ok(None),
};
let status_pid = status_file.as_ref().map(|status| status.pid);
let pid = pid_from_file.or(status_pid);
let live = pid_file_live || status_pid.map(process_exists).unwrap_or(false) || rpc_available;
if !live {
return Ok(None);
}
let target_root = normalized_abs_path(output_root)?;
let Some(daemon_root) = status_file
.as_ref()
.map(|status| status.resource_output_root.clone())
else {
return Ok(Some(DaemonResourceConflict {
pid,
daemon_root: None,
target_root,
}));
};
let normalized_daemon_root = normalized_abs_path(&daemon_root)?;
if normalized_daemon_root == target_root {
Ok(Some(DaemonResourceConflict {
pid,
daemon_root: Some(normalized_daemon_root),
target_root,
}))
} else {
Ok(None)
}
}
fn normalized_abs_path(path: &Path) -> anyhow::Result<PathBuf> {
let absolute = if path.is_absolute() {
path.to_path_buf()
} else {
env::current_dir()?.join(path)
};
Ok(fs::canonicalize(&absolute).unwrap_or_else(|_| lexically_normalize_path(&absolute)))
}
fn lexically_normalize_path(path: &Path) -> PathBuf {
let mut normalized = PathBuf::new();
for component in path.components() {
match component {
std::path::Component::CurDir => {}
std::path::Component::ParentDir => {
if !normalized.pop() {
normalized.push(component.as_os_str());
}
}
_ => normalized.push(component.as_os_str()),
}
}
if normalized.as_os_str().is_empty() {
PathBuf::from(".")
} else {
normalized
}
}
fn new_daemon_control() -> DaemonControl {
Arc::new((Mutex::new(DaemonControlState::default()), Condvar::new()))
}
fn daemon_control_stop_requested(control: Option<&DaemonControl>) -> bool {
let Some(control) = control else {
return false;
};
let (lock, _) = &**control;
lock.lock()
.map(|state| state.stop_requested)
.unwrap_or(true)
}
fn daemon_control_mark_stop_requested(control: &DaemonControl) {
let (lock, _) = &**control;
if let Ok(mut state) = lock.lock() {
state.stop_requested = true;
}
}
fn daemon_control_notify_all(control: &DaemonControl) {
let (_, cvar) = &**control;
cvar.notify_all();
}
fn daemon_control_request_refresh(control: &DaemonControl, force: bool) {
let (lock, cvar) = &**control;
if let Ok(mut state) = lock.lock() {
state.refresh_requested = true;
state.force_refresh_requested |= force;
}
cvar.notify_all();
}
fn daemon_control_request_reload(control: &DaemonControl) {
let (lock, cvar) = &**control;
if let Ok(mut state) = lock.lock() {
state.reload_requested = true;
}
cvar.notify_all();
}
fn wait_for_daemon_wake(control: Option<&DaemonControl>, timeout: Duration) -> DaemonWake {
let Some(control) = control else {
thread::sleep(timeout);
return DaemonWake::Timeout;
};
let (lock, cvar) = &**control;
let Ok(mut state) = lock.lock() else {
return DaemonWake::Stop;
};
if let Some(wake) = take_daemon_wake(&mut state) {
return wake;
}
let Ok((mut state, _timeout)) = cvar.wait_timeout(state, timeout) else {
return DaemonWake::Stop;
};
take_daemon_wake(&mut state).unwrap_or(DaemonWake::Timeout)
}
fn take_daemon_wake(state: &mut DaemonControlState) -> Option<DaemonWake> {
if state.stop_requested {
return Some(DaemonWake::Stop);
}
if state.reload_requested {
state.reload_requested = false;
state.refresh_requested = false;
state.force_refresh_requested = false;
return Some(DaemonWake::Reload);
}
if state.refresh_requested {
state.refresh_requested = false;
let force = state.force_refresh_requested;
state.force_refresh_requested = false;
return Some(DaemonWake::Refresh { force });
}
None
}
#[cfg(unix)]
fn start_daemon_rpc_server(
state_dir: &Path,
control: DaemonControl,
) -> anyhow::Result<DaemonRpcServer> {
fs::create_dir_all(state_dir)?;
let socket_path = daemon_socket_path(state_dir);
if socket_path.exists() {
match UnixStream::connect(&socket_path) {
Ok(_) => {
return Err(anyhow::anyhow!(
"后台 RPC socket 已被占用:{}",
socket_path.display()
));
}
Err(_) => {
let _ = fs::remove_file(&socket_path);
}
}
}
let listener = UnixListener::bind(&socket_path)?;
let server_state_dir = state_dir.to_path_buf();
thread::Builder::new()
.name("bat-daemon-rpc".to_string())
.spawn(move || {
for stream in listener.incoming() {
match stream {
Ok(stream) => {
let state_dir = server_state_dir.clone();
let control = Arc::clone(&control);
let _ = thread::Builder::new()
.name("bat-daemon-rpc-client".to_string())
.spawn(move || handle_daemon_rpc_client(stream, state_dir, control));
}
Err(error) => {
eprintln!("[daemon] RPC socket accept 失败:{error}");
break;
}
}
}
})?;
Ok(DaemonRpcServer { socket_path })
}
#[cfg(not(unix))]
fn start_daemon_rpc_server(
_state_dir: &Path,
_control: DaemonControl,
) -> anyhow::Result<DaemonRpcServer> {
Err(anyhow::anyhow!("daemon RPC 目前只支持 Unix/Linux 平台"))
}
#[cfg(unix)]
fn handle_daemon_rpc_client(mut stream: UnixStream, state_dir: PathBuf, control: DaemonControl) {
let Ok(reader_stream) = stream.try_clone() else {
return;
};
let mut reader = BufReader::new(reader_stream);
let mut line = String::new();
loop {
line.clear();
let bytes_read = match reader.read_line(&mut line) {
Ok(bytes_read) => bytes_read,
Err(error) => {
eprintln!("[daemon] RPC socket 读取失败:{error}");
return;
}
};
if bytes_read == 0 {
return;
}
if line.trim().is_empty() {
continue;
}
let mut notify_stop_after_response = false;
let response = match serde_json::from_str::<JsonRpcRequest>(&line) {
Ok(request) => {
notify_stop_after_response = request.method == RPC_METHOD_STOP;
handle_daemon_rpc_request(request, &state_dir, &control)
}
Err(error) => json_rpc_error(None, -32700, format!("JSON-RPC 请求解析失败:{error}")),
};
if let Err(error) = write_json_rpc_response(&mut stream, &response) {
if notify_stop_after_response {
daemon_control_notify_all(&control);
}
eprintln!("[daemon] RPC socket 写入失败:{error}");
return;
}
if notify_stop_after_response {
daemon_control_notify_all(&control);
}
}
}
fn handle_daemon_rpc_request(
request: JsonRpcRequest,
state_dir: &Path,
control: &DaemonControl,
) -> JsonRpcResponse {
let id = request.id.clone();
let result: anyhow::Result<serde_json::Value> = match request.method.as_str() {
RPC_METHOD_STATUS => {
let report = build_daemon_status_report(state_dir);
report.and_then(|report| serde_json::to_value(report).map_err(anyhow::Error::from))
}
RPC_METHOD_STOP => {
daemon_control_mark_stop_requested(control);
let _ = update_daemon_state_only(state_dir, "stopping");
serde_json::to_value(DaemonRpcAck {
command: "stop",
status: "accepted",
message: "后台停止请求已发送",
state_dir: state_dir.to_path_buf(),
socket_path: daemon_socket_path(state_dir),
force: None,
})
.map_err(anyhow::Error::from)
}
RPC_METHOD_RELOAD => {
daemon_control_request_reload(control);
serde_json::to_value(DaemonRpcAck {
command: "reload",
status: "accepted",
message: "后台重新加载请求已发送;空闲时会立即唤醒,忙碌时会在当前轮结束后重新执行自动发现和强制刷新",
state_dir: state_dir.to_path_buf(),
socket_path: daemon_socket_path(state_dir),
force: Some(true),
})
.map_err(anyhow::Error::from)
}
RPC_METHOD_REFRESH => {
let force = rpc_bool_param(request.params.as_ref(), "force").unwrap_or(false);
daemon_control_request_refresh(control, force);
serde_json::to_value(DaemonRpcAck {
command: "refresh",
status: "accepted",
message: if force {
"后台强制刷新请求已发送;空闲时会立即唤醒,忙碌时会在当前轮结束后执行"
} else {
"后台刷新检查请求已发送;空闲时会立即唤醒,忙碌时会在当前轮结束后执行"
},
state_dir: state_dir.to_path_buf(),
socket_path: daemon_socket_path(state_dir),
force: Some(force),
})
.map_err(anyhow::Error::from)
}
RPC_METHOD_LOGS => {
let tail = match rpc_tail_param(request.params.as_ref(), 200) {
Ok(tail) => tail,
Err(error) => return json_rpc_error(id, -32602, error.to_string()),
};
let report = build_logs_report(state_dir, tail);
report.and_then(|report| serde_json::to_value(report).map_err(anyhow::Error::from))
}
_ => {
return json_rpc_error(id, -32601, format!("未知后台 RPC 方法:{}", request.method));
}
};
match result {
Ok(value) => json_rpc_result(id, value),
Err(error) => json_rpc_error(id, -32603, error.to_string()),
}
}
#[cfg(unix)]
fn write_json_rpc_response(
stream: &mut UnixStream,
response: &JsonRpcResponse,
) -> anyhow::Result<()> {
serde_json::to_writer(&mut *stream, response)?;
stream.write_all(b"\n")?;
stream.flush()?;
Ok(())
}
fn json_rpc_result(id: Option<serde_json::Value>, result: serde_json::Value) -> JsonRpcResponse {
JsonRpcResponse {
jsonrpc: "2.0",
id,
result: Some(result),
error: None,
}
}
fn json_rpc_error(id: Option<serde_json::Value>, code: i32, message: String) -> JsonRpcResponse {
JsonRpcResponse {
jsonrpc: "2.0",
id,
result: None,
error: Some(JsonRpcError { code, message }),
}
}
fn rpc_bool_param(params: Option<&serde_json::Value>, key: &str) -> Option<bool> {
params
.and_then(|params| params.get(key))
.and_then(serde_json::Value::as_bool)
}
fn rpc_tail_param(params: Option<&serde_json::Value>, default: usize) -> anyhow::Result<usize> {
let tail = params
.and_then(|params| params.get("tail"))
.and_then(serde_json::Value::as_u64)
.unwrap_or(default as u64);
let tail = usize::try_from(tail)?;
if tail == 0 {
return Err(anyhow::anyhow!("tail 必须大于 0"));
}
Ok(tail)
}
#[cfg(unix)]
fn daemon_rpc_call(
state_dir: &Path,
method: &str,
params: Option<serde_json::Value>,
) -> anyhow::Result<serde_json::Value> {
let socket_path = daemon_socket_path(state_dir);
let mut stream = UnixStream::connect(&socket_path).map_err(|error| {
anyhow::anyhow!("无法连接后台 RPC socket {}{error}", socket_path.display())
})?;
let request = serde_json::json!({
"jsonrpc": "2.0",
"id": 1,
"method": method,
"params": params.unwrap_or(serde_json::Value::Null),
});
serde_json::to_writer(&mut stream, &request)?;
stream.write_all(b"\n")?;
stream.flush()?;
let mut reader = BufReader::new(stream);
let mut line = String::new();
reader.read_line(&mut line)?;
if line.trim().is_empty() {
return Err(anyhow::anyhow!("后台 RPC socket 未返回响应"));
}
let response: JsonRpcClientResponse = serde_json::from_str(&line)?;
if let Some(error) = response.error {
return Err(anyhow::anyhow!(
"后台 RPC {} 失败:{} ({})",
method,
error.message,
error.code
));
}
Ok(response.result.unwrap_or(serde_json::Value::Null))
}
#[cfg(not(unix))]
fn daemon_rpc_call(
_state_dir: &Path,
_method: &str,
_params: Option<serde_json::Value>,
) -> anyhow::Result<serde_json::Value> {
Err(anyhow::anyhow!("daemon RPC 目前只支持 Unix/Linux 平台"))
}
#[cfg(unix)]
fn daemon_rpc_available(state_dir: &Path) -> bool {
UnixStream::connect(daemon_socket_path(state_dir)).is_ok()
}
#[cfg(not(unix))]
fn daemon_rpc_available(_state_dir: &Path) -> bool {
false
}
#[derive(Debug, Serialize)]
struct DaemonStartReport {
status: &'static str,
message: &'static str,
pid: u32,
resource_output_root: PathBuf,
state_dir: PathBuf,
pid_path: PathBuf,
status_path: PathBuf,
log_path: PathBuf,
socket_path: PathBuf,
}
#[derive(Debug, Serialize)]
struct DaemonStatusReport {
status: &'static str,
message: &'static str,
running: bool,
pid: Option<u32>,
resource_output_root: Option<PathBuf>,
state_dir: PathBuf,
pid_path: PathBuf,
status_path: PathBuf,
socket_path: PathBuf,
rpc_available: bool,
stale_pid_file: bool,
stale_socket: bool,
log_path: Option<PathBuf>,
started_unix_seconds: Option<u64>,
updated_unix_seconds: Option<u64>,
daemon_state: Option<String>,
last_update_status: Option<String>,
last_error: Option<String>,
next_retry_seconds: Option<u64>,
pending_scheduled_force: Option<bool>,
next_forced_refresh_unix_seconds: Option<u64>,
command: Option<Vec<String>>,
}
#[derive(Debug, Serialize)]
struct DaemonStopReport {
status: &'static str,
message: &'static str,
stopped: bool,
pid: Option<u32>,
state_dir: PathBuf,
pid_path: PathBuf,
socket_path: PathBuf,
}
fn run_daemon_start(options: CliOptions) -> anyhow::Result<()> {
if options.config.dry_run {
return Err(anyhow::anyhow!("daemon 后台模式不能和 --dry-run 同时使用"));
}
if options.watch {
return Err(anyhow::anyhow!(
"--daemon 会自动启动 watch 模式,不要同时传 --watch"
));
}
let report = start_daemon_with_options(&options)?;
print_report(options.output_format, &report)?;
Ok(())
}
fn start_daemon_with_options(options: &CliOptions) -> anyhow::Result<DaemonStartReport> {
let resource_output_root = normalized_abs_path(&options.config.output_root)?;
let state_dir = options.state_dir.clone();
let args = daemon_child_args(options);
start_daemon_with_args(state_dir, resource_output_root, args, "后台同步已启动")
}
fn start_daemon_with_args(
state_dir: PathBuf,
resource_output_root: PathBuf,
args: Vec<String>,
message: &'static str,
) -> anyhow::Result<DaemonStartReport> {
fs::create_dir_all(&state_dir)?;
let pid_path = daemon_pid_path(&state_dir);
let status_path = daemon_status_path(&state_dir);
let log_path = daemon_log_path(&state_dir);
let socket_path = daemon_socket_path(&state_dir);
match classify_pid_lock_file(&pid_path)? {
PidLockState::Active(existing_pid) => {
return Err(anyhow::anyhow!(
"官方同步后台进程已经在运行,pid={existing_pid}"
));
}
PidLockState::StalePid(_) | PidLockState::Corrupt => {
let _ = fs::remove_file(&pid_path);
let _ = fs::remove_file(&socket_path);
}
PidLockState::Missing => {}
}
let executable = env::current_exe()?;
let log = OpenOptions::new()
.create(true)
.append(true)
.open(&log_path)?;
let log_for_stdout = log.try_clone()?;
let mut command = Command::new(&executable);
command
.args(&args)
.stdin(Stdio::null())
.stdout(Stdio::from(log_for_stdout))
.stderr(Stdio::from(log));
configure_daemon_command(&mut command);
let child = command.spawn()?;
let pid = child.id();
fs::write(&pid_path, pid.to_string())?;
let mut command = Vec::with_capacity(args.len() + 1);
command.push(executable.to_string_lossy().to_string());
command.extend(args);
let status = DaemonStatusFile {
version: DAEMON_STATUS_VERSION,
pid,
state: "started".to_string(),
resource_output_root: resource_output_root.clone(),
state_dir: state_dir.clone(),
log_path: log_path.clone(),
started_unix_seconds: unix_seconds_now(),
updated_unix_seconds: unix_seconds_now(),
last_update_status: None,
last_error: None,
next_retry_seconds: None,
pending_scheduled_force: false,
next_forced_refresh_unix_seconds: None,
command,
};
write_daemon_status_file(&status_path, &status)?;
Ok(DaemonStartReport {
status: "started",
message,
pid,
resource_output_root,
state_dir,
pid_path,
status_path,
log_path,
socket_path,
})
}
fn print_daemon_status(state_dir: &Path, output_format: OutputFormat) -> anyhow::Result<()> {
if let Ok(report) = daemon_rpc_call(state_dir, RPC_METHOD_STATUS, None) {
print_json_value(output_format, &report)?;
} else {
let report = build_daemon_status_report(state_dir)?;
print_report(output_format, &report)?;
}
Ok(())
}
fn stop_daemon(state_dir: &Path, output_format: OutputFormat) -> anyhow::Result<()> {
if daemon_rpc_available(state_dir) {
let pid = read_pid_file(&daemon_pid_path(state_dir))?.or_else(|| {
read_daemon_status_file(&daemon_status_path(state_dir))
.ok()
.flatten()
.map(|status| status.pid)
});
if daemon_rpc_call(state_dir, RPC_METHOD_STOP, None).is_ok() {
let forced_signal = if let Some(pid) = pid {
wait_for_rpc_stop_or_terminate(pid, Duration::from_secs(10))?
} else {
false
};
let pid_path = daemon_pid_path(state_dir);
let socket_path = daemon_socket_path(state_dir);
let _ = fs::remove_file(&pid_path);
let _ = fs::remove_file(&socket_path);
let report = DaemonStopReport {
status: if forced_signal {
"stopped_after_signal"
} else {
"stopped"
},
message: if forced_signal {
"后台进程已接收 RPC 停止请求,但未在优雅窗口内退出,已发送 SIGTERM 后停止"
} else {
"后台进程已通过 RPC 停止"
},
stopped: true,
pid,
state_dir: state_dir.to_path_buf(),
pid_path,
socket_path,
};
print_report(output_format, &report)?;
return Ok(());
}
}
let pid_path = daemon_pid_path(state_dir);
let Some(_pid) = read_pid_file(&pid_path)? else {
let report = stop_daemon_inner(state_dir)?;
print_report(output_format, &report)?;
return Ok(());
};
let report = stop_daemon_inner(state_dir)?;
print_report(output_format, &report)?;
Ok(())
}
fn stop_daemon_inner(state_dir: &Path) -> anyhow::Result<DaemonStopReport> {
let pid_path = daemon_pid_path(state_dir);
let pid_state = classify_pid_lock_file(&pid_path)?;
let (pid, running) = match pid_state {
PidLockState::Missing => {
return Ok(DaemonStopReport {
status: "not_running",
message: "后台进程当前未运行",
stopped: false,
pid: None,
state_dir: state_dir.to_path_buf(),
pid_path,
socket_path: daemon_socket_path(state_dir),
});
}
PidLockState::Corrupt => {
let _ = fs::remove_file(&pid_path);
let _ = fs::remove_file(daemon_socket_path(state_dir));
return Ok(DaemonStopReport {
status: "stale_pid_removed",
message: "已清理无效的后台 PID 文件",
stopped: false,
pid: None,
state_dir: state_dir.to_path_buf(),
pid_path,
socket_path: daemon_socket_path(state_dir),
});
}
PidLockState::StalePid(pid) => (pid, false),
PidLockState::Active(pid) => {
terminate_process(pid)?;
if !wait_for_process_exit(pid, Duration::from_secs(5)) {
return Err(anyhow::anyhow!("后台进程 pid={pid} 在 5 秒内未停止"));
}
(pid, true)
}
};
let _ = fs::remove_file(&pid_path);
let _ = fs::remove_file(daemon_socket_path(state_dir));
Ok(DaemonStopReport {
status: if running {
"stopped"
} else {
"stale_pid_removed"
},
message: if running {
"后台进程已停止"
} else {
"已清理失效的后台 PID 文件"
},
stopped: running,
pid: Some(pid),
state_dir: state_dir.to_path_buf(),
pid_path,
socket_path: daemon_socket_path(state_dir),
})
}
fn build_daemon_status_report(state_dir: &Path) -> anyhow::Result<DaemonStatusReport> {
let pid_path = daemon_pid_path(state_dir);
let status_path = daemon_status_path(state_dir);
let socket_path = daemon_socket_path(state_dir);
let status_file = read_daemon_status_file(&status_path)?;
let pid_file_state = classify_pid_lock_file(&pid_path)?;
let pid_file_pid = match pid_file_state {
PidLockState::Active(pid) | PidLockState::StalePid(pid) => Some(pid),
PidLockState::Missing | PidLockState::Corrupt => None,
};
let pid = pid_file_pid.or_else(|| status_file.as_ref().map(|status| status.pid));
let running = pid.map(process_exists).unwrap_or(false);
let rpc_available = daemon_rpc_available(state_dir);
let stale_pid_file = matches!(
pid_file_state,
PidLockState::StalePid(_) | PidLockState::Corrupt
);
let stale_socket = socket_path.exists() && !rpc_available;
Ok(DaemonStatusReport {
status: if running { "running" } else { "stopped" },
message: if running {
"后台进程正在运行"
} else if stale_pid_file || stale_socket {
"后台进程当前未运行,但检测到失效 PID/socket;可执行 clean-stable 清理"
} else {
"后台进程当前未运行"
},
running,
pid,
resource_output_root: status_file
.as_ref()
.map(|status| status.resource_output_root.clone()),
state_dir: state_dir.to_path_buf(),
pid_path,
status_path,
socket_path,
rpc_available,
stale_pid_file,
stale_socket,
log_path: status_file
.as_ref()
.map(|status| status.log_path.clone())
.or_else(|| Some(daemon_log_path(state_dir)).filter(|path| path.exists())),
started_unix_seconds: status_file
.as_ref()
.map(|status| status.started_unix_seconds),
updated_unix_seconds: status_file
.as_ref()
.map(|status| status.updated_unix_seconds),
daemon_state: status_file.as_ref().map(|status| status.state.clone()),
last_update_status: status_file
.as_ref()
.and_then(|status| status.last_update_status.clone()),
last_error: status_file
.as_ref()
.and_then(|status| status.last_error.clone()),
next_retry_seconds: status_file
.as_ref()
.and_then(|status| status.next_retry_seconds),
pending_scheduled_force: status_file
.as_ref()
.map(|status| status.pending_scheduled_force),
next_forced_refresh_unix_seconds: status_file
.as_ref()
.and_then(|status| status.next_forced_refresh_unix_seconds),
command: status_file.map(|status| status.command),
})
}
#[derive(Debug, Serialize)]
struct DaemonControlReport {
command: &'static str,
status: &'static str,
message: &'static str,
strategy: &'static str,
previous_pid: Option<u32>,
pid: u32,
state_dir: PathBuf,
resource_output_root: PathBuf,
log_path: PathBuf,
socket_path: PathBuf,
}
fn run_daemon_restart(options: &CliOptions, command_name: &'static str) -> anyhow::Result<()> {
let has_explicit_options = options.sync_option_explicit
|| options.output_explicit
|| tools_are_non_default(&options.config);
if command_name == "reload" && !has_explicit_options && daemon_rpc_available(&options.state_dir)
{
let report = daemon_rpc_call(&options.state_dir, RPC_METHOD_RELOAD, None)?;
print_json_value(options.output_format, &report)?;
return Ok(());
}
let status_file = read_daemon_status_file(&daemon_status_path(&options.state_dir))?;
let status_report = build_daemon_status_report(&options.state_dir)?;
let previous_pid = status_report.pid.filter(|pid| process_exists(*pid));
if previous_pid.is_some() {
if daemon_rpc_available(&options.state_dir)
&& daemon_rpc_call(&options.state_dir, RPC_METHOD_STOP, None).is_ok()
{
if let Some(pid) = previous_pid {
let _ = wait_for_rpc_stop_or_terminate(pid, Duration::from_secs(10))?;
}
let _ = fs::remove_file(daemon_pid_path(&options.state_dir));
let _ = fs::remove_file(daemon_socket_path(&options.state_dir));
} else {
let _ = stop_daemon_inner(&options.state_dir)?;
}
} else if read_pid_file(&daemon_pid_path(&options.state_dir))?.is_some() {
let _ = stop_daemon_inner(&options.state_dir)?;
}
let (state_dir, resource_output_root, args, strategy) = if has_explicit_options {
let mut start_options = options.clone();
start_options.daemon = true;
start_options.watch = false;
(
options.state_dir.clone(),
normalized_abs_path(&options.config.output_root)?,
daemon_child_args(&start_options),
"start_with_explicit_options",
)
} else {
let status_file = status_file.ok_or_else(|| {
anyhow::anyhow!(
"没有可复用的后台配置;请先执行 bat --auto-discover --daemon,或为 {command_name} 显式传入同步参数"
)
})?;
let mut command = status_file.command.into_iter();
let _executable = command.next().ok_or_else(|| {
anyhow::anyhow!("后台状态文件中的 command 为空,无法执行 {command_name}")
})?;
let args = command.collect::<Vec<_>>();
if args.is_empty() {
return Err(anyhow::anyhow!(
"后台状态文件中的 command 参数为空,无法执行 {command_name}"
));
}
(
options.state_dir.clone(),
status_file.resource_output_root,
args,
"restart_with_existing_command",
)
};
let start = start_daemon_with_args(
state_dir.clone(),
resource_output_root.clone(),
args,
if command_name == "reload" {
"后台配置已重新加载"
} else {
"后台进程已重启"
},
)?;
let report = DaemonControlReport {
command: command_name,
status: if command_name == "reload" {
"reloaded"
} else {
"restarted"
},
message: if command_name == "reload" {
"已使用原后台参数重启进程完成配置重新加载"
} else {
"已停止旧后台进程并启动新后台进程"
},
strategy,
previous_pid,
pid: start.pid,
state_dir,
resource_output_root,
log_path: start.log_path,
socket_path: start.socket_path,
};
print_report(options.output_format, &report)?;
Ok(())
}
#[derive(Debug, Serialize)]
struct CommandReport<T> {
command: &'static str,
status: &'static str,
message: &'static str,
data: T,
}
fn run_sync_command(options: &CliOptions, command_name: &'static str) -> anyhow::Result<()> {
if refresh_should_use_daemon_rpc(options, command_name)
&& daemon_rpc_available(&options.state_dir)
{
let _control_lock = DaemonControlLock::acquire(&options.state_dir)?;
let report = daemon_rpc_call(
&options.state_dir,
RPC_METHOD_REFRESH,
Some(serde_json::json!({ "force": options.config.force })),
)?;
print_json_value(options.output_format, &report)?;
return Ok(());
}
assert_no_live_daemon_output_conflict(options, command_name)?;
let mut config = options.config.clone();
if command_name == "repair" {
config.repair = true;
config.force = false;
}
let mut logger = ProgressLogger::new(options.progress);
let report =
OfficialUpdateService::new().run_with_progress(&config, |event| logger.log(event))?;
let command_report = CommandReport {
command: command_name,
status: "completed",
message: if command_name == "repair" {
"资源修复命令已执行"
} else if report.update_status == OfficialUpdateStatus::UpToDate {
"资源已是最新,无需更新"
} else {
"资源刷新命令已执行"
},
data: report,
};
print_report(options.output_format, &command_report)?;
Ok(())
}
fn refresh_should_use_daemon_rpc(options: &CliOptions, command_name: &str) -> bool {
if command_name != "refresh" {
return false;
}
let defaults = OfficialUpdateConfig::default();
options.command == CliCommand::Refresh
&& !options.watch
&& !options.daemon
&& !options.daemon_child
&& !options.output_explicit
&& options.config.server_info_source.is_none()
&& options.config.connection_group.is_none()
&& options.config.app_version.is_none()
&& options.config.launcher_version == defaults.launcher_version
&& options.config.platforms.is_none()
&& options.config.output_root == defaults.output_root
&& options.config.snapshot_path.is_none()
&& options.config.curl_command == defaults.curl_command
&& options.config.unzip_command == defaults.unzip_command
&& !options.config.dry_run
&& !options.config.plan
&& options.config.audit_local == defaults.audit_local
&& options.config.repair == defaults.repair
}
fn print_report<T>(format: OutputFormat, report: &T) -> anyhow::Result<()>
where
T: Serialize + HumanReport,
{
match format {
OutputFormat::Json => println!("{}", serde_json::to_string_pretty(report)?),
OutputFormat::Human => report.print_human()?,
}
Ok(())
}
fn print_json_value(format: OutputFormat, value: &serde_json::Value) -> anyhow::Result<()> {
match format {
OutputFormat::Json => println!("{}", serde_json::to_string_pretty(value)?),
OutputFormat::Human => print_human_json_value(value)?,
}
Ok(())
}
trait HumanReport {
fn print_human(&self) -> anyhow::Result<()>;
}
fn print_human_json_value(value: &serde_json::Value) -> anyhow::Result<()> {
if value.get("running").is_some() && value.get("state_dir").is_some() {
print_title("后台状态");
print_json_field(value, "status", "状态");
print_json_field(value, "message", "消息");
print_json_field(value, "running", "运行中");
print_json_field(value, "pid", "PID");
print_json_field(value, "daemon_state", "后台状态");
print_json_field(value, "rpc_available", "RPC 可用");
print_json_field(value, "stale_pid_file", "失效 PID");
print_json_field(value, "stale_socket", "失效 socket");
print_json_field(value, "last_update_status", "上次同步");
print_json_field(value, "last_error", "上次错误");
print_json_field(value, "next_retry_seconds", "下次重试秒数");
print_json_field(value, "resource_output_root", "资源目录");
print_json_field(value, "state_dir", "状态目录");
print_json_field(value, "socket_path", "socket");
print_json_field(value, "log_path", "日志");
return Ok(());
}
if value.get("command").and_then(serde_json::Value::as_str) == Some("logs") {
print_title("后台日志");
print_json_field(value, "status", "状态");
print_json_field(value, "message", "消息");
print_json_field(value, "log_path", "日志");
print_json_field(value, "bytes", "字节");
print_json_field(value, "total_lines", "总行数");
print_json_field(value, "returned_lines", "返回行数");
if let Some(content) = value.get("content").and_then(serde_json::Value::as_str) {
if !content.is_empty() {
println!();
println!("{content}");
}
}
return Ok(());
}
if value.get("command").is_some() && value.get("status").is_some() {
print_title("后台命令");
print_json_field(value, "command", "命令");
print_json_field(value, "status", "状态");
print_json_field(value, "message", "消息");
print_json_field(value, "force", "force");
print_json_field(value, "state_dir", "状态目录");
print_json_field(value, "socket_path", "socket");
return Ok(());
}
println!("{}", serde_json::to_string_pretty(value)?);
Ok(())
}
fn print_json_field(value: &serde_json::Value, key: &str, label: &str) {
let Some(value) = value.get(key) else {
return;
};
if value.is_null() {
return;
}
if let Some(value) = value.as_str() {
print_field(label, value);
} else {
print_field(label, value);
}
}
fn print_title(title: &str) {
println!("{title}");
}
fn print_field(label: &str, value: impl std::fmt::Display) {
println!(" {label:<18} {value}");
}
fn print_optional_field<T>(label: &str, value: Option<T>)
where
T: std::fmt::Display,
{
if let Some(value) = value {
print_field(label, value);
}
}
fn print_path_field(label: &str, value: &Path) {
print_field(label, value.display());
}
fn print_optional_path_field(label: &str, value: Option<&PathBuf>) {
if let Some(value) = value {
print_path_field(label, value);
}
}
fn print_list(label: &str, values: &[String], limit: usize) {
if values.is_empty() {
return;
}
println!(" {label}:");
for value in values.iter().take(limit) {
println!(" - {value}");
}
if values.len() > limit {
println!(" ... 还有 {} 项", values.len() - limit);
}
}
fn print_verification_summary(summary: &OfficialVerificationSummary) {
println!(" 校验摘要:");
for line in verification_summary_lines(summary) {
println!(" - {line}");
}
}
fn verification_summary_lines(summary: &OfficialVerificationSummary) -> Vec<String> {
vec![
format!(
"官方 .hash 强校验: {} 对 ({})",
summary.official_hash_verified_count, summary.official_hash_scope
),
format!(
"本地 BLAKE3 复用校验: {} 项通过, {} 项需修复 ({})",
summary.local_manifest_blake3_verified_count,
summary.local_manifest_repair_needed_count,
summary.local_manifest_blake3_scope
),
format!(
"ZIP 结构校验: {} 个 ZIP 通过 ({})",
summary.zip_structure_verified_count, summary.zip_structure_scope
),
]
}
fn format_bool(value: bool) -> &'static str {
if value {
"yes"
} else {
"no"
}
}
fn format_bytes(value: u64) -> String {
const KIB: f64 = 1024.0;
const MIB: f64 = 1024.0 * 1024.0;
const GIB: f64 = 1024.0 * 1024.0 * 1024.0;
let value_f = value as f64;
if value_f >= GIB {
format!("{value_f:.2} GiB", value_f = value_f / GIB)
} else if value_f >= MIB {
format!("{value_f:.2} MiB", value_f = value_f / MIB)
} else if value_f >= KIB {
format!("{value_f:.2} KiB", value_f = value_f / KIB)
} else {
format!("{value} B")
}
}
fn platform_label(platform: PatchPlatform) -> &'static str {
match platform {
PatchPlatform::Windows => "Windows",
PatchPlatform::Android => "Android",
}
}
impl HumanReport for OfficialUpdateReport {
fn print_human(&self) -> anyhow::Result<()> {
print_title("官方资源同步");
print_field("状态", self.update_status.as_str());
print_field("应用版本", &self.app_version);
print_optional_field("Bundle 版本", self.bundle_version.as_deref());
print_field("连接组", &self.connection_group);
print_field(
"平台",
self.platforms
.iter()
.map(|platform| platform_label(*platform))
.collect::<Vec<_>>()
.join(", "),
);
print_field("需要下载", format_bool(self.should_download));
print_field("首次同步", format_bool(self.is_initial));
print_field("强制刷新", format_bool(self.force));
print_field("本地审计", format_bool(self.audit_local));
print_field("自动修复", format_bool(self.repair));
print_field("dry-run", format_bool(self.dry_run));
print_path_field("snapshot", &self.snapshot_path);
print_path_field("manifest", &self.download_manifest);
print_optional_path_field("写入 snapshot", self.snapshot_written.as_ref());
print_optional_path_field("bootstrap cache", self.bootstrap_cache_path.as_ref());
print_optional_field("bootstrap 命中", self.bootstrap_cache_hit.map(format_bool));
print_optional_field("计划 URL 数", self.download_url_count);
print_optional_field("资源数", self.resource_count);
print_field("已下载", self.downloaded_count);
print_field("已续传", self.resumed_count);
print_field("已跳过", self.skipped_count);
print_field("传输量", format_bytes(self.transferred_bytes));
print_field("最终大小", format_bytes(self.final_bytes));
print_field("本地校验通过", self.local_manifest_verified_count);
print_field("需修复", self.local_manifest_repair_needed_count);
print_field("官方 hash 校验", self.official_seed_hash_verified_count);
print_verification_summary(&self.verification_summary);
print_field("catalog marker", self.addressables_marker_checked_count);
print_list("变更 endpoint", &self.changed_endpoint_urls, 8);
print_list("计划 URL", &self.download_urls, 8);
Ok(())
}
}
impl<T> HumanReport for CommandReport<T>
where
T: Serialize + HumanReport,
{
fn print_human(&self) -> anyhow::Result<()> {
print_title(self.message);
print_field("命令", self.command);
print_field("状态", self.status);
self.data.print_human()
}
}
impl HumanReport for DaemonStartReport {
fn print_human(&self) -> anyhow::Result<()> {
print_title(self.message);
print_field("状态", self.status);
print_field("PID", self.pid);
print_path_field("资源目录", &self.resource_output_root);
print_path_field("状态目录", &self.state_dir);
print_path_field("socket", &self.socket_path);
print_path_field("日志", &self.log_path);
Ok(())
}
}
impl HumanReport for DaemonStatusReport {
fn print_human(&self) -> anyhow::Result<()> {
print_title("后台状态");
print_field("状态", self.status);
print_field("消息", self.message);
print_field("运行中", format_bool(self.running));
print_optional_field("PID", self.pid);
print_optional_field("后台状态", self.daemon_state.as_deref());
print_field("RPC 可用", format_bool(self.rpc_available));
print_field("失效 PID", format_bool(self.stale_pid_file));
print_field("失效 socket", format_bool(self.stale_socket));
print_optional_field("上次同步", self.last_update_status.as_deref());
print_optional_field("上次错误", self.last_error.as_deref());
print_optional_field("下次重试秒数", self.next_retry_seconds);
print_optional_path_field("资源目录", self.resource_output_root.as_ref());
print_path_field("状态目录", &self.state_dir);
print_path_field("socket", &self.socket_path);
print_optional_path_field("日志", self.log_path.as_ref());
Ok(())
}
}
impl HumanReport for DaemonStopReport {
fn print_human(&self) -> anyhow::Result<()> {
print_title(self.message);
print_field("状态", self.status);
print_field("已停止", format_bool(self.stopped));
print_optional_field("PID", self.pid);
print_path_field("状态目录", &self.state_dir);
print_path_field("socket", &self.socket_path);
Ok(())
}
}
impl HumanReport for DaemonControlReport {
fn print_human(&self) -> anyhow::Result<()> {
print_title(self.message);
print_field("命令", self.command);
print_field("状态", self.status);
print_field("策略", self.strategy);
print_optional_field("旧 PID", self.previous_pid);
print_field("PID", self.pid);
print_path_field("资源目录", &self.resource_output_root);
print_path_field("状态目录", &self.state_dir);
print_path_field("socket", &self.socket_path);
print_path_field("日志", &self.log_path);
Ok(())
}
}
impl HumanReport for VerifyCommandReport {
fn print_human(&self) -> anyhow::Result<()> {
print_title(self.message);
print_field("命令", self.command);
print_field("状态", self.status);
print_field("健康", format_bool(self.healthy));
print_field("远端状态", &self.remote_update_status);
print_optional_field("计划 URL 数", self.planned_url_count);
print_field("计划异常数", self.expected_plan_failure_count);
print_field("本地 manifest 项", self.local_manifest_entry_count);
print_field("本地校验通过", self.local_manifest_verified_count);
print_field("本地失败数", self.local_manifest_failure_count);
print_field("官方 hash 对", self.official_hash_pair_count);
print_field("官方 hash 通过", self.official_hash_verified_count);
print_verification_summary(&self.verification_summary);
if !self.official_hash_errors.is_empty() {
print_list("官方 hash 错误", &self.official_hash_errors, 8);
}
if !self.failures.is_empty() {
println!(" 失败项:");
for item in self.failures.iter().take(12) {
println!(" - {} -> {}", item.status, item.destination.display());
}
if self.failures.len() > 12 {
println!(" ... 还有 {} 项", self.failures.len() - 12);
}
}
Ok(())
}
}
impl HumanReport for LogsReport {
fn print_human(&self) -> anyhow::Result<()> {
print_title(self.message);
print_field("状态", self.status);
print_path_field("日志", &self.log_path);
print_field("存在", format_bool(self.exists));
print_field("为空", format_bool(self.empty));
print_field("字节", self.bytes);
print_field("总行数", self.total_lines);
print_field("返回行数", self.returned_lines);
if !self.content.is_empty() {
println!();
println!("{}", self.content);
}
Ok(())
}
}
impl HumanReport for DoctorReport {
fn print_human(&self) -> anyhow::Result<()> {
print_title(self.message);
print_field("命令", self.command);
print_field("状态", self.status);
print_field("健康", format_bool(self.healthy));
println!(" 检查:");
for check in &self.checks {
println!(
" [{}] {} - {}",
if check.ok { "OK" } else { "FAIL" },
check.name,
check.message
);
}
Ok(())
}
}
impl HumanReport for CleanStableReport {
fn print_human(&self) -> anyhow::Result<()> {
print_title(self.message);
print_field("命令", self.command);
print_field("状态", self.status);
print_path_field("资源目录", &self.output_root);
print_path_field("状态目录", &self.state_dir);
if !self.removed_paths.is_empty() {
println!(" 已清理:");
for path in &self.removed_paths {
println!(" - {}", path.display());
}
}
if !self.skipped_paths.is_empty() {
println!(" 已跳过:");
for path in &self.skipped_paths {
println!(" - {}", path.display());
}
}
Ok(())
}
}
impl HumanReport for DaemonRpcAck {
fn print_human(&self) -> anyhow::Result<()> {
print_title(self.message);
print_field("命令", self.command);
print_field("状态", self.status);
print_optional_field("force", self.force.map(format_bool));
print_path_field("状态目录", &self.state_dir);
print_path_field("socket", &self.socket_path);
Ok(())
}
}
#[derive(Debug, Serialize)]
struct VerificationItemReport {
url: String,
destination: PathBuf,
status: String,
expected_bytes: Option<u64>,
actual_bytes: Option<u64>,
expected_blake3: Option<String>,
actual_blake3: Option<String>,
zip_error: Option<String>,
}
#[derive(Debug, Serialize)]
struct VerifyCommandReport {
command: &'static str,
status: &'static str,
message: &'static str,
healthy: bool,
remote_update_status: String,
planned_url_count: Option<usize>,
expected_plan_failure_count: usize,
local_manifest_entry_count: usize,
local_manifest_verified_count: usize,
local_manifest_failure_count: usize,
official_hash_pair_count: usize,
official_hash_verified_count: usize,
verification_summary: OfficialVerificationSummary,
official_hash_errors: Vec<String>,
failures: Vec<VerificationItemReport>,
}
fn run_verify_command(options: &CliOptions) -> anyhow::Result<bool> {
let mut config = options.config.clone();
config.dry_run = true;
config.plan = true;
config.audit_local = true;
config.repair = false;
config.force = false;
let mut logger = ProgressLogger::new(options.progress);
let update_report =
OfficialUpdateService::new().run_with_progress(&config, |event| logger.log(event))?;
let verification = bat_infrastructure::OfficialResourcePullService::with_curl_command(
&config.output_root,
&config.curl_command,
)
.verify_local_download_manifest()
.map_err(anyhow::Error::msg)?;
let failures = verification
.items
.iter()
.filter(|item| !item.status.is_verified())
.map(|item| VerificationItemReport {
url: item.url.clone(),
destination: item.destination.clone(),
status: item.status.as_str().to_string(),
expected_bytes: item.expected_bytes,
actual_bytes: item.actual_bytes,
expected_blake3: item.expected_blake3.clone(),
actual_blake3: item.actual_blake3.clone(),
zip_error: item.zip_error.clone(),
})
.collect::<Vec<_>>();
let healthy = update_report.update_status == OfficialUpdateStatus::UpToDate
&& update_report.local_manifest_repair_needed_count == 0
&& verification.is_clean();
let verification_summary = OfficialVerificationSummary::new(
verification.manifest_blake3_verified_count(),
verification.failure_count(),
verification.official_hash_verified_count,
verification.zip_structure_verified_count(),
);
let report = VerifyCommandReport {
command: "verify",
status: if healthy { "verified" } else { "failed" },
message: if healthy {
"所有当前计划资源和本地 manifest 均通过验证"
} else {
"发现资源或官方 seed hash 校验异常"
},
healthy,
remote_update_status: update_report.update_status.as_str().to_string(),
planned_url_count: update_report.download_url_count,
expected_plan_failure_count: update_report.local_manifest_repair_needed_count,
local_manifest_entry_count: verification.items.len(),
local_manifest_verified_count: verification.verified_count(),
local_manifest_failure_count: verification.failure_count(),
official_hash_pair_count: verification.official_hash_pair_count,
official_hash_verified_count: verification.official_hash_verified_count,
verification_summary,
official_hash_errors: verification.official_hash_errors,
failures,
};
print_report(options.output_format, &report)?;
Ok(healthy)
}
#[derive(Debug, Serialize)]
struct LogsReport {
command: &'static str,
status: &'static str,
message: &'static str,
log_path: PathBuf,
exists: bool,
empty: bool,
bytes: usize,
total_lines: usize,
returned_lines: usize,
content: String,
}
fn run_logs_command(options: &CliOptions) -> anyhow::Result<()> {
if daemon_rpc_available(&options.state_dir) {
let report = daemon_rpc_call(
&options.state_dir,
RPC_METHOD_LOGS,
Some(serde_json::json!({ "tail": options.tail_lines })),
)?;
print_json_value(options.output_format, &report)?;
return Ok(());
}
let report = build_logs_report(&options.state_dir, options.tail_lines)?;
print_report(options.output_format, &report)?;
Ok(())
}
fn build_logs_report(state_dir: &Path, tail_lines: usize) -> anyhow::Result<LogsReport> {
let status_file = read_daemon_status_file(&daemon_status_path(state_dir))?;
let log_path = status_file
.as_ref()
.map(|status| status.log_path.clone())
.unwrap_or_else(|| daemon_log_path(state_dir));
let mut exists = true;
let bytes = match fs::read(&log_path) {
Ok(bytes) => bytes,
Err(error) if error.kind() == std::io::ErrorKind::NotFound => Vec::new(),
Err(error) => return Err(error.into()),
};
if bytes.is_empty() && !log_path.exists() {
exists = false;
}
let raw = String::from_utf8_lossy(&bytes);
let lines = raw.lines().collect::<Vec<_>>();
let start = lines.len().saturating_sub(tail_lines);
let content = lines[start..].join("\n");
Ok(LogsReport {
command: "logs",
status: if bytes.is_empty() { "empty" } else { "ok" },
message: if bytes.is_empty() {
"后台日志为空或尚未创建"
} else {
"已读取后台日志"
},
log_path,
exists,
empty: bytes.is_empty(),
bytes: bytes.len(),
total_lines: lines.len(),
returned_lines: lines.len().saturating_sub(start),
content,
})
}
#[derive(Debug, Serialize)]
struct DoctorCheck {
name: &'static str,
ok: bool,
message: String,
}
#[derive(Debug, Serialize)]
struct DoctorReport {
command: &'static str,
status: &'static str,
message: &'static str,
healthy: bool,
checks: Vec<DoctorCheck>,
}
fn run_doctor_command(options: &CliOptions) -> anyhow::Result<bool> {
let mut checks = Vec::new();
checks.push(path_check(
"state_dir",
&options.state_dir,
"后台状态目录可用",
));
checks.push(path_check(
"output_root",
&options.config.output_root,
"资源输出目录可用",
));
checks.push(command_check("curl", &options.config.curl_command));
checks.push(command_check("unzip", &options.config.unzip_command));
let pid_path = daemon_pid_path(&options.state_dir);
let daemon_running = read_pid_file(&pid_path)
.ok()
.flatten()
.is_some_and(process_exists);
match read_pid_file(&pid_path) {
Ok(Some(pid)) if process_exists(pid) => checks.push(DoctorCheck {
name: "daemon",
ok: true,
message: format!("后台进程正在运行,pid={pid}"),
}),
Ok(Some(pid)) => checks.push(DoctorCheck {
name: "daemon",
ok: true,
message: format!("发现失效 PID 文件,pid={pid};可执行 clean-stable 清理"),
}),
Ok(None) => checks.push(DoctorCheck {
name: "daemon",
ok: true,
message: "后台进程当前未运行".to_string(),
}),
Err(error) => checks.push(DoctorCheck {
name: "daemon",
ok: false,
message: format!("读取后台 PID 失败:{error}"),
}),
}
let socket_path = daemon_socket_path(&options.state_dir);
let socket_exists = socket_path.exists();
let socket_available = daemon_rpc_available(&options.state_dir);
checks.push(DoctorCheck {
name: "daemon_rpc",
ok: if daemon_running {
socket_available
} else {
true
},
message: if socket_available {
format!("后台 RPC socket 可连接:{}", socket_path.display())
} else if daemon_running {
format!(
"后台进程正在运行,但 RPC socket 不可连接:{}",
socket_path.display()
)
} else if socket_exists {
format!(
"发现失效后台 RPC socket{};可执行 clean-stable 清理",
socket_path.display()
)
} else {
"后台 RPC socket 尚未创建".to_string()
},
});
let lock_path = options.config.lock_path();
checks.push(match classify_pid_lock_file(&lock_path)? {
PidLockState::Missing => DoctorCheck {
name: "resource_lock",
ok: true,
message: "资源目录没有活动锁".to_string(),
},
PidLockState::Active(pid) => DoctorCheck {
name: "resource_lock",
ok: false,
message: format!(
"资源目录锁正在使用:{}owner_pid={pid}",
lock_path.display()
),
},
PidLockState::StalePid(pid) => DoctorCheck {
name: "resource_lock",
ok: false,
message: format!(
"发现失效资源锁:{}owner_pid={pid};可执行 clean-stable 清理",
lock_path.display()
),
},
PidLockState::Corrupt => DoctorCheck {
name: "resource_lock",
ok: false,
message: format!(
"发现损坏资源锁:{};可执行 clean-stable 清理",
lock_path.display()
),
},
});
let control_lock_path = daemon_control_lock_path(&options.state_dir);
checks.push(match classify_pid_lock_file(&control_lock_path)? {
PidLockState::Missing => DoctorCheck {
name: "daemon_control_lock",
ok: true,
message: "后台控制锁没有活动锁".to_string(),
},
PidLockState::Active(pid) => DoctorCheck {
name: "daemon_control_lock",
ok: true,
message: format!(
"后台控制命令正在运行:{}owner_pid={pid}",
control_lock_path.display()
),
},
PidLockState::StalePid(pid) => DoctorCheck {
name: "daemon_control_lock",
ok: false,
message: format!(
"发现失效后台控制锁:{}owner_pid={pid};可执行 clean-stable 清理",
control_lock_path.display()
),
},
PidLockState::Corrupt => DoctorCheck {
name: "daemon_control_lock",
ok: false,
message: format!(
"发现损坏后台控制锁:{};可执行 clean-stable 清理",
control_lock_path.display()
),
},
});
let status_path = daemon_status_path(&options.state_dir);
checks.push(DoctorCheck {
name: "daemon_status",
ok: match read_daemon_status_file(&status_path) {
Ok(Some(_)) | Ok(None) => true,
Err(_) => false,
},
message: if status_path.exists() {
format!("后台状态文件可解析:{}", status_path.display())
} else {
"后台状态文件尚未创建".to_string()
},
});
let healthy = checks.iter().all(|check| check.ok);
let report = DoctorReport {
command: "doctor",
status: if healthy { "ok" } else { "issues_found" },
message: if healthy {
"运行时诊断通过"
} else {
"运行时诊断发现问题"
},
healthy,
checks,
};
print_report(options.output_format, &report)?;
Ok(healthy)
}
fn path_check(name: &'static str, path: &Path, ready_message: &str) -> DoctorCheck {
let (ok, message) = if path.is_dir() {
(true, format!("{ready_message}{}", path.display()))
} else if path.exists() {
(false, format!("路径存在但不是目录:{}", path.display()))
} else {
let parent = path.parent().unwrap_or_else(|| Path::new("."));
(
parent.is_dir(),
if parent.is_dir() {
format!(
"目录尚未创建,但父目录可用,将在首次运行时创建:{}",
path.display()
)
} else {
format!("目录不存在且父目录不可用:{}", path.display())
},
)
};
DoctorCheck { name, ok, message }
}
fn command_check(name: &'static str, command: &Path) -> DoctorCheck {
let result = Command::new(command).arg("--version").output();
DoctorCheck {
name,
ok: result.is_ok(),
message: match result {
Ok(output) => format!(
"命令可执行:{}(退出码={:?}",
command.display(),
output.status.code()
),
Err(error) => format!("命令不可执行 {}{error}", command.display()),
},
}
}
#[derive(Debug, Serialize)]
struct CleanStableReport {
command: &'static str,
status: &'static str,
message: &'static str,
output_root: PathBuf,
state_dir: PathBuf,
removed_paths: Vec<PathBuf>,
skipped_paths: Vec<PathBuf>,
}
fn run_clean_stable_command(options: &CliOptions) -> anyhow::Result<()> {
let status = build_daemon_status_report(&options.state_dir)?;
if status.running || status.rpc_available {
return Err(anyhow::anyhow!(
"后台进程正在运行,不能执行 clean-stable;请先执行 stop"
));
}
let mut removed_paths = Vec::new();
let mut skipped_paths = Vec::new();
for root in [&options.config.output_root, &options.state_dir] {
collect_transient_files(root, &mut removed_paths)?;
}
let pid_path = daemon_pid_path(&options.state_dir);
match classify_pid_lock_file(&pid_path)? {
PidLockState::Missing => {}
PidLockState::Active(_) => skipped_paths.push(pid_path),
PidLockState::StalePid(_) | PidLockState::Corrupt => {
fs::remove_file(&pid_path)?;
removed_paths.push(pid_path);
}
}
let socket_path = daemon_socket_path(&options.state_dir);
if socket_path.exists() {
if daemon_rpc_available(&options.state_dir) {
skipped_paths.push(socket_path);
} else {
fs::remove_file(&socket_path)?;
removed_paths.push(socket_path);
}
}
let lock_path = options.config.lock_path();
match classify_pid_lock_file(&lock_path)? {
PidLockState::Missing => {}
PidLockState::Active(_) => skipped_paths.push(lock_path),
PidLockState::StalePid(_) | PidLockState::Corrupt => {
fs::remove_file(&lock_path)?;
removed_paths.push(lock_path);
}
}
let control_lock_path = daemon_control_lock_path(&options.state_dir);
match classify_pid_lock_file(&control_lock_path)? {
PidLockState::Missing => {}
PidLockState::Active(pid) => {
return Err(anyhow::anyhow!(
"后台控制命令已被锁定 (locked){}owner_pid={pid} 仍在运行",
control_lock_path.display()
));
}
PidLockState::StalePid(_) | PidLockState::Corrupt => {
fs::remove_file(&control_lock_path)?;
removed_paths.push(control_lock_path);
}
}
let report = CleanStableReport {
command: "clean-stable",
status: if skipped_paths.is_empty() {
"cleaned"
} else {
"cleaned_with_skips"
},
message: "已清理断点临时文件、临时状态文件和失效锁/PID;不会删除正式资源",
output_root: options.config.output_root.clone(),
state_dir: options.state_dir.clone(),
removed_paths,
skipped_paths,
};
print_report(options.output_format, &report)?;
Ok(())
}
fn collect_transient_files(root: &Path, removed_paths: &mut Vec<PathBuf>) -> anyhow::Result<()> {
if !root.exists() {
return Ok(());
}
let metadata = fs::symlink_metadata(root)?;
if !metadata.is_dir() {
return Ok(());
}
for entry in fs::read_dir(root)? {
let entry = entry?;
let path = entry.path();
let metadata = fs::symlink_metadata(&path)?;
if metadata.is_dir() {
collect_transient_files(&path, removed_paths)?;
continue;
}
let file_name = path
.file_name()
.and_then(|name| name.to_str())
.unwrap_or("");
if file_name.ends_with(".part") || file_name.ends_with(".tmp") {
fs::remove_file(&path)?;
removed_paths.push(path);
}
}
Ok(())
}
fn read_daemon_status_file(path: &Path) -> anyhow::Result<Option<DaemonStatusFile>> {
let bytes = match fs::read(path) {
Ok(bytes) => bytes,
Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(None),
Err(error) => return Err(error.into()),
};
Ok(Some(serde_json::from_slice(&bytes)?))
}
fn write_daemon_status_file(path: &Path, status: &DaemonStatusFile) -> anyhow::Result<()> {
if let Some(parent) = path.parent() {
fs::create_dir_all(parent)?;
}
let temporary = path.with_extension("json.tmp");
fs::write(&temporary, serde_json::to_vec_pretty(status)?)?;
fs::rename(&temporary, path)?;
Ok(())
}
fn update_daemon_status(
state_dir: &Path,
state: &str,
last_update_status: Option<String>,
last_error: Option<String>,
next_retry_seconds: Option<u64>,
pending_scheduled_force: bool,
next_forced_refresh_at: SystemTime,
) -> anyhow::Result<()> {
let status_path = daemon_status_path(state_dir);
let mut status = read_daemon_status_file(&status_path)?.unwrap_or_else(|| DaemonStatusFile {
version: DAEMON_STATUS_VERSION,
pid: std::process::id(),
state: "started".to_string(),
resource_output_root: OfficialUpdateConfig::default().output_root,
state_dir: state_dir.to_path_buf(),
log_path: daemon_log_path(state_dir),
started_unix_seconds: unix_seconds_now(),
updated_unix_seconds: unix_seconds_now(),
last_update_status: None,
last_error: None,
next_retry_seconds: None,
pending_scheduled_force: false,
next_forced_refresh_unix_seconds: None,
command: env::args().collect(),
});
status.pid = std::process::id();
status.state = state.to_string();
status.updated_unix_seconds = unix_seconds_now();
status.last_update_status = last_update_status;
status.last_error = last_error;
status.next_retry_seconds = next_retry_seconds;
status.pending_scheduled_force = pending_scheduled_force;
status.next_forced_refresh_unix_seconds = system_time_to_unix_seconds(next_forced_refresh_at);
write_daemon_status_file(&status_path, &status)
}
fn update_daemon_state_only(state_dir: &Path, state: &str) -> anyhow::Result<()> {
let status_path = daemon_status_path(state_dir);
let Some(mut status) = read_daemon_status_file(&status_path)? else {
return Ok(());
};
status.pid = std::process::id();
status.state = state.to_string();
status.updated_unix_seconds = unix_seconds_now();
status.next_retry_seconds = None;
write_daemon_status_file(&status_path, &status)
}
fn read_pid_file(path: &Path) -> anyhow::Result<Option<u32>> {
let contents = match fs::read_to_string(path) {
Ok(contents) => contents,
Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(None),
Err(error) => return Err(error.into()),
};
Ok(contents.trim().parse::<u32>().ok().filter(|pid| *pid > 0))
}
fn daemon_pid_path(output_root: &Path) -> PathBuf {
output_root.join(DAEMON_PID_FILE)
}
fn daemon_status_path(output_root: &Path) -> PathBuf {
output_root.join(DAEMON_STATUS_FILE)
}
fn daemon_log_path(output_root: &Path) -> PathBuf {
output_root.join(DAEMON_LOG_FILE)
}
fn daemon_socket_path(output_root: &Path) -> PathBuf {
output_root.join(DAEMON_SOCKET_FILE)
}
fn daemon_control_lock_path(output_root: &Path) -> PathBuf {
output_root.join(DAEMON_CONTROL_LOCK_FILE)
}
fn daemon_child_args(options: &CliOptions) -> Vec<String> {
let mut args = Vec::new();
let config = &options.config;
if config.auto_discover {
args.push("--auto-discover".to_string());
}
if let Some(source) = config.server_info_source.as_ref() {
match source {
OfficialServerInfoSource::LocalPath(path) => {
args.push("--server-info-path".to_string());
args.push(path.to_string_lossy().to_string());
}
OfficialServerInfoSource::OfficialFile(name) => {
args.push("--server-info-file".to_string());
args.push(name.clone());
}
OfficialServerInfoSource::OfficialUrl(url) => {
args.push("--server-info-url".to_string());
args.push(url.clone());
}
}
}
if let Some(connection_group) = config.connection_group.as_ref() {
args.push("--connection-group".to_string());
args.push(connection_group.clone());
}
if let Some(app_version) = config.app_version.as_ref() {
args.push("--app-version".to_string());
args.push(app_version.clone());
}
args.push("--launcher-version".to_string());
args.push(config.launcher_version.clone());
if let Some(platforms) = config.platforms.as_ref() {
args.push("--platforms".to_string());
args.push(
platforms
.iter()
.map(|platform| match platform {
PatchPlatform::Windows => "Windows",
PatchPlatform::Android => "Android",
})
.collect::<Vec<_>>()
.join(","),
);
}
args.push("--output".to_string());
args.push(config.output_root.to_string_lossy().to_string());
if let Some(snapshot_path) = config.snapshot_path.as_ref() {
args.push("--snapshot".to_string());
args.push(snapshot_path.to_string_lossy().to_string());
}
args.push("--curl".to_string());
args.push(config.curl_command.to_string_lossy().to_string());
args.push("--unzip".to_string());
args.push(config.unzip_command.to_string_lossy().to_string());
if config.force {
args.push("--force".to_string());
}
if config.audit_local {
args.push("--audit-local".to_string());
} else {
args.push("--no-audit-local".to_string());
}
if config.repair {
args.push("--repair".to_string());
} else {
args.push("--no-repair".to_string());
}
args.push("--state-dir".to_string());
args.push(options.state_dir.to_string_lossy().to_string());
args.push("--daemon-child".to_string());
args.push("--watch".to_string());
args.push("--interval".to_string());
args.push(format_duration_arg(options.interval));
args.push("--error-retry".to_string());
args.push(format_duration_arg(options.error_retry_interval));
if options.quiet_up_to_date {
args.push("--quiet-up-to-date".to_string());
} else {
args.push("--no-quiet-up-to-date".to_string());
}
if options.progress {
args.push("--progress".to_string());
} else {
args.push("--no-progress".to_string());
}
if options.output_format == OutputFormat::Json {
args.push("--json".to_string());
}
if options.banner {
args.push("--banner".to_string());
} else {
args.push("--no-banner".to_string());
}
args
}
fn format_duration_arg(duration: Duration) -> String {
let millis = duration.as_millis();
if millis % 3_600_000 == 0 {
format!("{}h", millis / 3_600_000)
} else if millis % 60_000 == 0 {
format!("{}m", millis / 60_000)
} else if millis % 1_000 == 0 {
format!("{}s", millis / 1_000)
} else {
format!("{millis}ms")
}
}
fn unix_seconds_now() -> u64 {
SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap_or_default()
.as_secs()
}
fn system_time_to_unix_seconds(time: SystemTime) -> Option<u64> {
time.duration_since(UNIX_EPOCH)
.ok()
.map(|duration| duration.as_secs())
}
#[cfg(unix)]
fn configure_daemon_command(command: &mut Command) {
use std::os::unix::process::CommandExt;
unsafe {
command.pre_exec(|| {
if libc::setsid() < 0 {
Err(std::io::Error::last_os_error())
} else {
Ok(())
}
});
}
}
#[cfg(not(unix))]
fn configure_daemon_command(_command: &mut Command) {}
#[cfg(unix)]
fn process_exists(pid: u32) -> bool {
let Ok(pid) = libc::pid_t::try_from(pid) else {
return false;
};
if pid <= 0 {
return false;
}
unsafe {
libc::kill(pid, 0) == 0
|| std::io::Error::last_os_error().raw_os_error() == Some(libc::EPERM)
}
}
#[cfg(not(unix))]
fn process_exists(_pid: u32) -> bool {
true
}
#[cfg(unix)]
fn terminate_process(pid: u32) -> anyhow::Result<()> {
let pid = libc::pid_t::try_from(pid)?;
if pid <= 0 {
return Err(anyhow::anyhow!("无效的后台进程 pid"));
}
let rc = unsafe { libc::kill(pid, libc::SIGTERM) };
if rc == 0 {
Ok(())
} else {
Err(std::io::Error::last_os_error().into())
}
}
#[cfg(not(unix))]
fn terminate_process(_pid: u32) -> anyhow::Result<()> {
Err(anyhow::anyhow!("stop 目前只支持 Unix/Linux 平台"))
}
fn wait_for_process_exit(pid: u32, timeout: Duration) -> bool {
let started_at = Instant::now();
while started_at.elapsed() < timeout {
if !process_exists(pid) {
return true;
}
thread::sleep(Duration::from_millis(100));
}
!process_exists(pid)
}
fn wait_for_rpc_stop_or_terminate(pid: u32, timeout: Duration) -> anyhow::Result<bool> {
if wait_for_process_exit(pid, timeout) {
return Ok(false);
}
terminate_process(pid)?;
if !wait_for_process_exit(pid, Duration::from_secs(5)) {
return Err(anyhow::anyhow!("后台进程 pid={pid} 在 SIGTERM 后仍未停止"));
}
Ok(true)
}
fn next_forced_refresh_at_or_after(now: SystemTime) -> SystemTime {
forced_refresh_time(now, ForcedRefreshBoundary::AtOrAfter)
}
fn next_forced_refresh_after(now: SystemTime) -> SystemTime {
forced_refresh_time(now, ForcedRefreshBoundary::After)
}
#[derive(Debug, Clone, Copy, PartialEq, Eq)]
enum ForcedRefreshBoundary {
AtOrAfter,
After,
}
fn forced_refresh_time(now: SystemTime, boundary: ForcedRefreshBoundary) -> SystemTime {
let local_seconds = beijing_local_seconds_since_epoch(now);
let local_day = local_seconds / SECONDS_PER_DAY;
let local_second_of_day = local_seconds % SECONDS_PER_DAY;
for scheduled_second in DAILY_FORCED_REFRESH_LOCAL_SECONDS {
let matches_boundary = match boundary {
ForcedRefreshBoundary::AtOrAfter => scheduled_second >= local_second_of_day,
ForcedRefreshBoundary::After => scheduled_second > local_second_of_day,
};
if matches_boundary {
return system_time_from_beijing_local_seconds(
local_day
.saturating_mul(SECONDS_PER_DAY)
.saturating_add(scheduled_second),
);
}
}
system_time_from_beijing_local_seconds(
local_day
.saturating_add(1)
.saturating_mul(SECONDS_PER_DAY)
.saturating_add(DAILY_FORCED_REFRESH_LOCAL_SECONDS[0]),
)
}
fn beijing_local_seconds_since_epoch(time: SystemTime) -> u64 {
time.duration_since(UNIX_EPOCH)
.unwrap_or_default()
.as_secs()
.saturating_add(BEIJING_UTC_OFFSET_SECONDS)
}
fn system_time_from_beijing_local_seconds(local_seconds: u64) -> SystemTime {
UNIX_EPOCH + Duration::from_secs(local_seconds.saturating_sub(BEIJING_UTC_OFFSET_SECONDS))
}
fn duration_until(deadline: SystemTime, now: SystemTime) -> Duration {
deadline.duration_since(now).unwrap_or_default()
}
#[derive(Debug, Clone)]
struct ProgressLogger {
enabled: bool,
started_at: Instant,
}
impl ProgressLogger {
fn new(enabled: bool) -> Self {
Self {
enabled,
started_at: Instant::now(),
}
}
fn log(&mut self, event: OfficialUpdateProgress) {
self.log_text(event.stage, event.message);
}
fn log_text(&mut self, stage: &str, message: impl AsRef<str>) {
if !self.enabled {
return;
}
eprintln!(
"[+{} 信息] [{}] {}",
format_duration(self.started_at.elapsed()),
localized_stage(stage),
message.as_ref()
);
}
}
fn localized_stage(stage: &str) -> &str {
match stage {
"start" => "启动",
"lock" => "锁",
"bootstrap" => "启动发现",
"launcher" => "启动器",
"bootstrap-cache" => "启动缓存",
"game-main-config" => "游戏配置",
"metadata" => "元数据",
"server-info" => "服务器信息",
"discovery" => "发现",
"markers" => "标记",
"marker" => "标记",
"snapshot" => "快照",
"decision" => "决策",
"plan" => "计划",
"catalog" => "目录",
"inventory" => "清单",
"local-state" => "本地状态",
"audit" => "审计",
"dry-run" => "试运行",
"download" => "下载",
"resource" => "资源",
"finish" => "完成",
"watch" => "常驻",
"daemon" => "后台",
_ => stage,
}
}
fn is_locked_error(message: &str) -> bool {
message.contains("locked") || message.contains("锁定")
}
fn localize_error_message(message: &str) -> String {
message.to_string()
}
fn print_startup_banner() {
eprintln!("{STARTUP_BANNER}");
}
fn should_print_status(status: OfficialUpdateStatus, quiet_up_to_date: bool) -> bool {
!(quiet_up_to_date && status == OfficialUpdateStatus::UpToDate)
}
fn parse_args() -> anyhow::Result<CliOptions> {
parse_args_from(env::args())
}
fn parse_args_from(raw_args: impl IntoIterator<Item = String>) -> anyhow::Result<CliOptions> {
let mut args = raw_args.into_iter();
let binary = args.next().unwrap_or_else(|| "bat".to_string());
let mut options = CliOptions::default();
while let Some(flag) = args.next() {
match flag.as_str() {
"status" => {
ensure_command_not_set(options.command, "status")?;
options.command = CliCommand::Status;
}
"stop" => {
ensure_command_not_set(options.command, "stop")?;
options.command = CliCommand::Stop;
}
"restart" => {
ensure_command_not_set(options.command, "restart")?;
options.command = CliCommand::Restart;
}
"reload" => {
ensure_command_not_set(options.command, "reload")?;
options.command = CliCommand::Reload;
}
"refresh" => {
ensure_command_not_set(options.command, "refresh")?;
options.command = CliCommand::Refresh;
}
"verify" => {
ensure_command_not_set(options.command, "verify")?;
options.command = CliCommand::Verify;
}
"repair" => {
ensure_command_not_set(options.command, "repair")?;
options.command = CliCommand::Repair;
}
"doctor" => {
ensure_command_not_set(options.command, "doctor")?;
options.command = CliCommand::Doctor;
}
"logs" => {
ensure_command_not_set(options.command, "logs")?;
options.command = CliCommand::Logs;
}
"clean-stable" => {
ensure_command_not_set(options.command, "clean-stable")?;
options.command = CliCommand::CleanStable;
}
"--auto-discover" => {
options.config.auto_discover = true;
options.sync_option_explicit = true;
}
"--server-info-url" => {
options.config.server_info_source = Some(OfficialServerInfoSource::OfficialUrl(
next_option_value(&mut args, &flag)?,
));
options.sync_option_explicit = true;
}
"--server-info-file" => {
options.config.server_info_source = Some(OfficialServerInfoSource::OfficialFile(
next_option_value(&mut args, &flag)?,
));
options.sync_option_explicit = true;
}
"--server-info-path" => {
options.config.server_info_source = Some(OfficialServerInfoSource::LocalPath(
PathBuf::from(next_option_value(&mut args, &flag)?),
));
options.sync_option_explicit = true;
}
"--connection-group" => {
options.config.connection_group = Some(next_option_value(&mut args, &flag)?);
options.sync_option_explicit = true;
}
"--app-version" => {
options.config.app_version = Some(next_option_value(&mut args, &flag)?);
options.sync_option_explicit = true;
}
"--launcher-version" => {
options.config.launcher_version = next_option_value(&mut args, &flag)?;
options.sync_option_explicit = true;
}
"--platforms" => {
let value = next_option_value(&mut args, &flag)?;
options.config.platforms =
Some(parse_platforms(&value).map_err(anyhow::Error::msg)?);
options.sync_option_explicit = true;
}
"--output" => {
options.config.output_root = PathBuf::from(next_option_value(&mut args, &flag)?);
options.output_explicit = true;
}
"--state-dir" | "--pid-dir" => {
options.state_dir = PathBuf::from(next_option_value(&mut args, &flag)?);
}
"--snapshot" => {
options.config.snapshot_path =
Some(PathBuf::from(next_option_value(&mut args, &flag)?));
options.sync_option_explicit = true;
}
"--curl" => {
options.config.curl_command = PathBuf::from(next_option_value(&mut args, &flag)?);
}
"--unzip" => {
options.config.unzip_command = PathBuf::from(next_option_value(&mut args, &flag)?);
}
"--dry-run" => {
options.config.dry_run = true;
options.sync_option_explicit = true;
}
"--plan" => {
options.config.plan = true;
options.sync_option_explicit = true;
}
"--force" => {
options.config.force = true;
options.sync_option_explicit = true;
}
"--audit-local" => {
options.config.audit_local = true;
options.sync_option_explicit = true;
}
"--no-audit-local" => {
options.config.audit_local = false;
options.sync_option_explicit = true;
}
"--repair" => {
options.config.repair = true;
options.sync_option_explicit = true;
}
"--no-repair" => {
options.config.repair = false;
options.sync_option_explicit = true;
}
"--watch" => {
options.watch = true;
options.sync_option_explicit = true;
}
"--daemon" => {
options.daemon = true;
options.sync_option_explicit = true;
}
"--daemon-child" => {
options.daemon_child = true;
options.watch = true;
options.sync_option_explicit = true;
}
"--interval" => {
options.interval = parse_duration(&next_option_value(&mut args, &flag)?)?;
options.sync_option_explicit = true;
}
"--interval-seconds" => {
let seconds = next_option_value(&mut args, &flag)?
.parse::<u64>()
.map_err(|error| anyhow::anyhow!("{flag} 的秒数无效:{error}"))?;
options.interval = Duration::from_secs(seconds);
options.sync_option_explicit = true;
}
"--error-retry" => {
options.error_retry_interval =
parse_duration(&next_option_value(&mut args, &flag)?)?;
options.sync_option_explicit = true;
}
"--error-retry-seconds" => {
let seconds = next_option_value(&mut args, &flag)?
.parse::<u64>()
.map_err(|error| anyhow::anyhow!("{flag} 的秒数无效:{error}"))?;
options.error_retry_interval = Duration::from_secs(seconds);
options.sync_option_explicit = true;
}
"--quiet-up-to-date" => {
options.quiet_up_to_date = true;
options.quiet_up_to_date_explicit = true;
options.sync_option_explicit = true;
}
"--no-quiet-up-to-date" => {
options.quiet_up_to_date = false;
options.quiet_up_to_date_explicit = true;
options.sync_option_explicit = true;
}
"--progress" => {
options.progress = true;
options.banner = true;
}
"--no-progress" => {
options.progress = false;
options.banner = false;
}
"--json" => {
options.output_format = OutputFormat::Json;
}
"--human" => {
options.output_format = OutputFormat::Human;
}
"--banner" => {
options.banner = true;
}
"--no-banner" => {
options.banner = false;
}
"--tail" => {
options.tail_lines = next_option_value(&mut args, &flag)?
.parse::<usize>()
.map_err(|error| anyhow::anyhow!("{flag} 的行数无效:{error}"))?;
if options.tail_lines == 0 {
return Err(anyhow::anyhow!("--tail 必须大于 0"));
}
}
"--help" | "-h" => {
print_usage(&binary);
std::process::exit(0);
}
_ => return Err(anyhow::anyhow!("未知参数:{flag}")),
}
}
match options.command {
CliCommand::Status | CliCommand::Stop | CliCommand::Logs => {
if options.sync_option_explicit
|| options.output_explicit
|| tools_are_non_default(&options.config)
{
return Err(anyhow::anyhow!(
"status/stop/logs 只读取 --state-dir;资源同步参数和输出目录无关"
));
}
options.progress = false;
options.banner = false;
}
CliCommand::Doctor | CliCommand::CleanStable => {
if options.sync_option_explicit {
return Err(anyhow::anyhow!(
"doctor/clean-stable 只接受 --output、--state-dir、--curl 和 --unzip 等诊断参数"
));
}
options.progress = false;
options.banner = false;
}
CliCommand::Refresh | CliCommand::Verify | CliCommand::Repair => {
if !options.config.auto_discover
&& options.config.server_info_source.is_none()
&& options.config.connection_group.is_none()
&& options.config.app_version.is_none()
{
options.config.auto_discover = true;
}
if matches!(options.command, CliCommand::Verify) {
options.config.dry_run = true;
options.config.plan = true;
options.config.audit_local = true;
options.config.repair = false;
}
}
CliCommand::Restart | CliCommand::Reload => {
if (options.sync_option_explicit
|| options.output_explicit
|| tools_are_non_default(&options.config))
&& !options.config.auto_discover
&& options.config.server_info_source.is_none()
&& options.config.connection_group.is_none()
&& options.config.app_version.is_none()
{
options.config.auto_discover = true;
}
}
CliCommand::Run => {}
}
if options.daemon && options.watch {
return Err(anyhow::anyhow!(
"--daemon 会自动启动 watch 模式,不要同时传 --watch"
));
}
if (options.watch || options.daemon || options.daemon_child) && options.config.dry_run {
return Err(anyhow::anyhow!(
"watch/daemon 模式不能和 --dry-run 同时使用"
));
}
if options.interval.is_zero() {
return Err(anyhow::anyhow!("watch 检查间隔必须大于 0"));
}
if options.error_retry_interval.is_zero() {
return Err(anyhow::anyhow!("watch 失败重试间隔必须大于 0"));
}
if (options.watch || options.daemon || options.daemon_child)
&& !options.quiet_up_to_date_explicit
{
options.quiet_up_to_date = true;
}
if matches!(options.command, CliCommand::Verify) {
if options.config.force {
return Err(anyhow::anyhow!(
"verify 只验证本地资源,不能和 --force 同时使用"
));
}
if options.config.repair == false && options.config.audit_local == false {
return Err(anyhow::anyhow!(
"verify 必须启用本地审计;不要同时传 --no-repair --no-audit-local"
));
}
}
if matches!(options.command, CliCommand::Restart | CliCommand::Reload)
&& (options.config.force || options.config.dry_run)
{
return Err(anyhow::anyhow!(
"restart/reload 不能和 --force 或 --dry-run 同时使用"
));
}
if matches!(options.command, CliCommand::Logs) && options.tail_lines == 0 {
return Err(anyhow::anyhow!("logs 的 --tail 必须大于 0"));
}
Ok(options)
}
fn ensure_command_not_set(command: CliCommand, next: &str) -> anyhow::Result<()> {
if command == CliCommand::Run {
Ok(())
} else {
Err(anyhow::anyhow!("不支持同时指定多个命令;意外命令:{next}"))
}
}
fn tools_are_non_default(config: &OfficialUpdateConfig) -> bool {
config.curl_command != PathBuf::from("curl") || config.unzip_command != PathBuf::from("unzip")
}
fn next_option_value(
args: &mut impl Iterator<Item = String>,
flag: &str,
) -> anyhow::Result<String> {
args.next()
.ok_or_else(|| anyhow::anyhow!("{flag} 缺少参数值"))
}
fn print_usage(binary: &str) {
eprintln!("BlueArchiveToolkit official resource sync");
eprintln!();
eprintln!("Usage:");
eprintln!(" {binary} [OPTIONS]");
eprintln!(" {binary} <COMMAND> [OPTIONS]");
eprintln!();
eprintln!("Commands:");
eprintln!(" refresh Run one update check, or ask a live daemon to refresh");
eprintln!(" verify Verify remote plan, local manifest, and official seed hashes");
eprintln!(" repair Redownload resources that fail local verification");
eprintln!(" status Show daemon state");
eprintln!(" stop Stop daemon");
eprintln!(
" restart Restart daemon, reusing saved args unless explicit args are passed"
);
eprintln!(" reload Ask daemon to rediscover metadata and force refresh");
eprintln!(" logs Show daemon log tail");
eprintln!(" doctor Run runtime diagnostics");
eprintln!(" clean-stable Remove .part/.tmp/stale lock, pid, and socket files");
eprintln!();
eprintln!("Examples:");
eprintln!(" {binary} --auto-discover --dry-run");
eprintln!(" {binary} --auto-discover --watch");
eprintln!(" {binary} --auto-discover --daemon");
eprintln!(" {binary} status");
eprintln!(" {binary} refresh --force --json");
eprintln!();
eprintln!("Discovery:");
eprintln!(
" --auto-discover Discover app-version, server-info, connection-group"
);
eprintln!(" --server-info-url <URL> Use an official server-info URL");
eprintln!(" --server-info-file <NAME> Use an official server-info file name");
eprintln!(" --server-info-path <PATH> Use a local server-info JSON file");
eprintln!(" --app-version <VERSION> Override app version");
eprintln!(" --connection-group <NAME> Override connection group");
eprintln!(" --launcher-version <VERSION> Launcher metadata API version (default: 1.7.2)");
eprintln!();
eprintln!("Sync:");
eprintln!(" --platforms <LIST> Platforms, e.g. Windows,Android");
eprintln!(" --output <DIR> Resource output dir (default: ./bat-resources)");
eprintln!(" --snapshot <PATH> Snapshot path (default: <output>/official-sync-snapshot.json)");
eprintln!(" --curl <PATH> curl executable (default: curl)");
eprintln!(" --unzip <PATH> unzip executable (default: unzip)");
eprintln!(" --dry-run Do not write sync state");
eprintln!(" --plan Include planned URLs in dry-run");
eprintln!(" --force Force download/refresh");
eprintln!(" --audit-local | --no-audit-local Enable/disable local manifest audit");
eprintln!(" --repair | --no-repair Enable/disable automatic repair");
eprintln!();
eprintln!("Daemon:");
eprintln!(" --watch Run in foreground loop");
eprintln!(" --daemon Start detached watch process");
eprintln!(" --state-dir <DIR> Daemon state dir (default: /tmp/bat-pid)");
eprintln!(" --interval <DURATION> Normal check interval (default: 1h)");
eprintln!(" --error-retry <DURATION> Retry interval after error (default: 60s)");
eprintln!(" --quiet-up-to-date Suppress clean up-to-date reports");
eprintln!(" --no-quiet-up-to-date Always print reports");
eprintln!(" --tail <N> Log lines for logs command (default: 200)");
eprintln!();
eprintln!("Output:");
eprintln!(" --human Human-readable output (default)");
eprintln!(" --json Stable JSON output for scripts");
eprintln!(" --progress | --no-progress Enable/disable stderr progress logs");
eprintln!(" --banner | --no-banner Enable/disable startup banner");
eprintln!(" -h, --help Show this help");
eprintln!();
eprintln!("Defaults:");
eprintln!(" platforms: Windows,Android");
eprintln!(" resource output: ./bat-resources");
eprintln!(" daemon state: /tmp/bat-pid (bat.sock, bat.pid, bat-status.json, bat-daemon.log)");
eprintln!(" forced refresh: {DAILY_FORCED_REFRESH_LABEL}");
}
fn parse_duration(value: &str) -> anyhow::Result<Duration> {
let value = value.trim();
if value.is_empty() {
return Err(anyhow::anyhow!("时间间隔不能为空"));
}
let (number, multiplier) = if let Some(number) = value.strip_suffix("ms") {
let millis = number
.parse::<u64>()
.map_err(|error| anyhow::anyhow!("毫秒时间间隔无效:{value}{error}"))?;
return Ok(Duration::from_millis(millis));
} else if let Some(number) = value.strip_suffix('s') {
(number, 1)
} else if let Some(number) = value.strip_suffix('m') {
(number, 60)
} else if let Some(number) = value.strip_suffix('h') {
(number, 60 * 60)
} else {
(value, 1)
};
let amount = number
.parse::<u64>()
.map_err(|error| anyhow::anyhow!("时间间隔无效:{value}{error}"))?;
Ok(Duration::from_secs(amount.saturating_mul(multiplier)))
}
fn format_duration(duration: Duration) -> String {
let total_millis = duration.as_millis();
let minutes = total_millis / 60_000;
let seconds = (total_millis / 1_000) % 60;
let millis = total_millis % 1_000;
format!("{minutes:02}:{seconds:02}.{millis:03}")
}
fn parse_platforms(value: &str) -> Result<Vec<PatchPlatform>, String> {
value
.split(',')
.map(str::trim)
.filter(|part| !part.is_empty())
.map(parse_platform)
.collect()
}
fn parse_platform(value: &str) -> Result<PatchPlatform, String> {
match value.to_ascii_lowercase().as_str() {
"windows" | "win" => Ok(PatchPlatform::Windows),
"android" => Ok(PatchPlatform::Android),
_ => Err(format!("不支持的平台:{value}")),
}
}
#[cfg(test)]
mod tests {
use super::*;
fn parse(values: &[&str]) -> anyhow::Result<CliOptions> {
parse_args_from(values.iter().map(|value| value.to_string()))
}
#[test]
fn parses_auto_discover_sync_args() {
let options = parse(&[
"bat",
"--auto-discover",
"--platforms",
"Windows,Android",
"--snapshot",
"/tmp/snapshot.json",
"--dry-run",
"--plan",
])
.unwrap();
let config = options.config;
assert!(config.auto_discover);
assert_eq!(
config.platforms.as_deref(),
Some([PatchPlatform::Windows, PatchPlatform::Android].as_slice())
);
assert_eq!(
config.snapshot_path,
Some(PathBuf::from("/tmp/snapshot.json"))
);
assert!(config.dry_run);
assert!(config.plan);
assert!(config.audit_local);
assert!(config.repair);
assert!(!options.watch);
assert_eq!(
options.interval,
Duration::from_secs(DEFAULT_WATCH_INTERVAL_SECONDS)
);
assert_eq!(
options.error_retry_interval,
Duration::from_secs(DEFAULT_ERROR_RETRY_SECONDS)
);
assert!(!options.quiet_up_to_date);
assert!(!options.quiet_up_to_date_explicit);
assert!(options.progress);
assert!(options.banner);
assert_eq!(options.output_format, OutputFormat::Human);
assert_eq!(config.output_root, PathBuf::from("./bat-resources"));
assert_eq!(options.state_dir, PathBuf::from(DEFAULT_DAEMON_STATE_DIR));
}
#[test]
fn parses_explicit_source_and_disable_repair() {
let options = parse(&[
"bat",
"--server-info-url",
"https://yostar-serverinfo.bluearchiveyostar.com/r93_fixture.json",
"--connection-group",
"Prod-Audit",
"--app-version",
"1.70.0",
"--no-audit-local",
"--no-repair",
"--unzip",
"/usr/bin/unzip",
])
.unwrap();
let config = options.config;
assert!(!config.auto_discover);
assert_eq!(config.connection_group.as_deref(), Some("Prod-Audit"));
assert_eq!(config.app_version.as_deref(), Some("1.70.0"));
assert!(!config.audit_local);
assert!(!config.repair);
assert_eq!(config.unzip_command, PathBuf::from("/usr/bin/unzip"));
assert!(matches!(
config.server_info_source,
Some(OfficialServerInfoSource::OfficialUrl(ref url))
if url == "https://yostar-serverinfo.bluearchiveyostar.com/r93_fixture.json"
));
}
#[test]
fn parses_watch_defaults_to_one_hour_and_quiet_up_to_date() {
let options = parse(&["bat", "--auto-discover", "--watch", "--interval", "30m"]).unwrap();
assert!(options.watch);
assert_eq!(options.interval, Duration::from_secs(30 * 60));
assert_eq!(
options.error_retry_interval,
Duration::from_secs(DEFAULT_ERROR_RETRY_SECONDS)
);
assert!(options.quiet_up_to_date);
assert!(!options.quiet_up_to_date_explicit);
}
#[test]
fn parses_daemon_mode_as_background_watch() {
let options = parse(&[
"bat",
"--auto-discover",
"--daemon",
"--output",
"/tmp/daemon-output",
"--interval",
"30m",
])
.unwrap();
assert_eq!(options.command, CliCommand::Run);
assert!(options.daemon);
assert!(!options.watch);
assert_eq!(
options.config.output_root,
PathBuf::from("/tmp/daemon-output")
);
assert_eq!(options.state_dir, PathBuf::from(DEFAULT_DAEMON_STATE_DIR));
assert_eq!(options.interval, Duration::from_secs(30 * 60));
assert!(options.quiet_up_to_date);
}
#[test]
fn parses_output_format_flags() {
let options = parse(&["bat", "status", "--json"]).unwrap();
assert_eq!(options.command, CliCommand::Status);
assert_eq!(options.output_format, OutputFormat::Json);
let options = parse(&["bat", "status", "--json", "--human"]).unwrap();
assert_eq!(options.output_format, OutputFormat::Human);
}
#[test]
fn parses_status_and_stop_commands() {
let status = parse(&["bat", "status"]).unwrap();
assert_eq!(status.command, CliCommand::Status);
assert!(!status.watch);
assert!(!status.daemon);
assert!(!status.progress);
assert!(!status.banner);
assert_eq!(status.config.output_root, PathBuf::from("./bat-resources"));
assert_eq!(status.state_dir, PathBuf::from(DEFAULT_DAEMON_STATE_DIR));
let stop = parse(&["bat", "stop", "--state-dir", "/tmp/custom-bat-pid"]).unwrap();
assert_eq!(stop.command, CliCommand::Stop);
assert_eq!(stop.config.output_root, PathBuf::from("./bat-resources"));
assert_eq!(stop.state_dir, PathBuf::from("/tmp/custom-bat-pid"));
}
#[test]
fn parses_management_and_resource_commands() {
let refresh = parse(&["bat", "refresh", "--force"]).unwrap();
assert_eq!(refresh.command, CliCommand::Refresh);
assert!(refresh.config.auto_discover);
assert!(refresh.config.force);
let verify = parse(&["bat", "verify"]).unwrap();
assert_eq!(verify.command, CliCommand::Verify);
assert!(verify.config.auto_discover);
assert!(verify.config.dry_run);
assert!(verify.config.plan);
assert!(!verify.config.repair);
let repair = parse(&["bat", "repair"]).unwrap();
assert_eq!(repair.command, CliCommand::Repair);
assert!(repair.config.auto_discover);
for (value, expected) in [
("restart", CliCommand::Restart),
("reload", CliCommand::Reload),
("logs", CliCommand::Logs),
("doctor", CliCommand::Doctor),
("clean-stable", CliCommand::CleanStable),
] {
let options = parse(&["bat", value]).unwrap();
assert_eq!(options.command, expected);
}
}
#[test]
fn parses_logs_tail_and_rejects_force_for_verify() {
let options = parse(&["bat", "logs", "--tail", "25"]).unwrap();
assert_eq!(options.tail_lines, 25);
let error = parse(&["bat", "verify", "--force"]).unwrap_err();
assert!(error.to_string().contains("不能和 --force"));
}
#[test]
fn restart_with_explicit_options_defaults_to_auto_discover() {
let options = parse(&["bat", "restart", "--output", "/tmp/bat-resources"]).unwrap();
assert_eq!(options.command, CliCommand::Restart);
assert!(options.output_explicit);
assert!(options.config.auto_discover);
}
#[test]
fn daemon_child_args_preserve_sync_options() {
let options = parse(&[
"bat",
"--auto-discover",
"--daemon",
"--platforms",
"Windows,Android",
"--output",
"/tmp/daemon-output",
"--state-dir",
"/tmp/daemon-state",
"--curl",
"/usr/bin/curl",
"--unzip",
"/usr/bin/unzip",
"--interval",
"30m",
"--error-retry",
"45s",
"--json",
"--no-banner",
])
.unwrap();
let args = daemon_child_args(&options);
assert!(args.contains(&"--watch".to_string()));
assert!(args.contains(&"--daemon-child".to_string()));
assert!(!args.contains(&"--daemon".to_string()));
assert!(args
.windows(2)
.any(|pair| pair == ["--output", "/tmp/daemon-output"]));
assert!(args
.windows(2)
.any(|pair| pair == ["--state-dir", "/tmp/daemon-state"]));
assert!(args
.windows(2)
.any(|pair| pair == ["--platforms", "Windows,Android"]));
assert!(args.windows(2).any(|pair| pair == ["--interval", "30m"]));
assert!(args.windows(2).any(|pair| pair == ["--error-retry", "45s"]));
assert!(args.contains(&"--quiet-up-to-date".to_string()));
assert!(args.contains(&"--json".to_string()));
assert!(args.contains(&"--no-banner".to_string()));
}
#[test]
fn parses_explicit_watch_error_retry_interval() {
let options =
parse(&["bat", "--watch", "--interval", "1h", "--error-retry", "45s"]).unwrap();
assert_eq!(options.interval, Duration::from_secs(3600));
assert_eq!(options.error_retry_interval, Duration::from_secs(45));
let options = parse(&["bat", "--watch", "--error-retry-seconds", "30"]).unwrap();
assert_eq!(options.error_retry_interval, Duration::from_secs(30));
}
#[test]
fn parses_progress_flags() {
let options = parse(&["bat"]).unwrap();
assert!(options.progress);
assert!(options.banner);
let options = parse(&["bat", "--no-progress"]).unwrap();
assert!(!options.progress);
assert!(!options.banner);
let options = parse(&["bat", "--no-progress", "--progress"]).unwrap();
assert!(options.progress);
assert!(options.banner);
}
#[test]
fn parses_banner_flags() {
let options = parse(&["bat", "--no-banner"]).unwrap();
assert!(options.progress);
assert!(!options.banner);
let options = parse(&["bat", "--no-progress", "--banner"]).unwrap();
assert!(!options.progress);
assert!(options.banner);
}
#[test]
fn startup_banner_contains_product_name() {
assert!(STARTUP_BANNER.contains("BlueArchiveToolkit"));
}
#[test]
fn verification_summary_output_names_all_validation_types() {
let summary = OfficialVerificationSummary::new(7, 1, 2, 3);
let lines = verification_summary_lines(&summary);
assert_eq!(lines.len(), 3);
assert!(lines[0].contains("官方 .hash 强校验"));
assert!(lines[0].contains("xxHash32"));
assert!(lines[1].contains("本地 BLAKE3 复用校验"));
assert!(lines[1].contains("需修复"));
assert!(lines[2].contains("ZIP 结构校验"));
assert!(lines[2].contains("central directory"));
let json = serde_json::to_value(&summary).unwrap();
assert_eq!(json["official_hash_verified_count"], 2);
assert_eq!(json["local_manifest_blake3_verified_count"], 7);
assert_eq!(json["local_manifest_repair_needed_count"], 1);
assert_eq!(json["zip_structure_verified_count"], 3);
}
#[test]
fn daemon_socket_path_uses_state_dir() {
assert_eq!(
daemon_socket_path(Path::new("/tmp/custom-bat-pid")),
PathBuf::from("/tmp/custom-bat-pid/bat.sock")
);
}
#[test]
fn rpc_tail_param_defaults_and_rejects_zero() {
assert_eq!(rpc_tail_param(None, 200).unwrap(), 200);
assert_eq!(
rpc_tail_param(Some(&serde_json::json!({ "tail": 25 })), 200).unwrap(),
25
);
let error = rpc_tail_param(Some(&serde_json::json!({ "tail": 0 })), 200).unwrap_err();
assert!(error.to_string().contains("大于 0"));
}
#[test]
fn json_rpc_request_and_response_are_line_protocol_friendly() {
let request: JsonRpcRequest = serde_json::from_value(serde_json::json!({
"jsonrpc": "2.0",
"id": 7,
"method": "bat.refresh",
"params": { "force": true }
}))
.unwrap();
assert_eq!(request.method, RPC_METHOD_REFRESH);
assert_eq!(rpc_bool_param(request.params.as_ref(), "force"), Some(true));
let response = json_rpc_result(
Some(serde_json::json!(7)),
serde_json::json!({ "ok": true }),
);
let serialized = serde_json::to_string(&response).unwrap();
assert!(serialized.contains("\"jsonrpc\":\"2.0\""));
assert!(serialized.contains("\"ok\":true"));
assert!(!serialized.contains('\n'));
}
#[test]
fn daemon_control_wake_handles_refresh_reload_and_stop() {
let control = new_daemon_control();
daemon_control_request_refresh(&control, true);
assert_eq!(
wait_for_daemon_wake(Some(&control), Duration::from_millis(1)),
DaemonWake::Refresh { force: true }
);
daemon_control_request_reload(&control);
assert_eq!(
wait_for_daemon_wake(Some(&control), Duration::from_millis(1)),
DaemonWake::Reload
);
daemon_control_request_refresh(&control, false);
daemon_control_mark_stop_requested(&control);
daemon_control_notify_all(&control);
assert_eq!(
wait_for_daemon_wake(Some(&control), Duration::from_millis(1)),
DaemonWake::Stop
);
}
#[test]
fn daemon_control_lock_rejects_active_owner() {
let temp = tempfile::TempDir::new().unwrap();
let state_dir = temp.path().join("state");
let first = DaemonControlLock::acquire(&state_dir).unwrap();
let second = DaemonControlLock::acquire(&state_dir).unwrap_err();
assert!(second.to_string().contains("locked"));
drop(first);
assert!(DaemonControlLock::acquire(&state_dir).is_ok());
}
#[test]
fn daemon_control_lock_recovers_stale_and_corrupt_files() {
let temp = tempfile::TempDir::new().unwrap();
let state_dir = temp.path().join("state");
fs::create_dir_all(&state_dir).unwrap();
let lock_path = daemon_control_lock_path(&state_dir);
fs::write(&lock_path, "999999999").unwrap();
let lock = DaemonControlLock::acquire(&state_dir).unwrap();
drop(lock);
assert!(!lock_path.exists());
fs::write(&lock_path, "not-a-pid").unwrap();
let lock = DaemonControlLock::acquire(&state_dir).unwrap();
assert_eq!(
fs::read_to_string(&lock_path).unwrap(),
std::process::id().to_string()
);
drop(lock);
assert!(!lock_path.exists());
}
#[test]
fn live_daemon_resource_conflict_detects_same_output() {
let temp = tempfile::TempDir::new().unwrap();
let state_dir = temp.path().join("state");
let output_root = temp.path().join("resources");
fs::create_dir_all(&output_root).unwrap();
write_daemon_status_file(
&daemon_status_path(&state_dir),
&test_daemon_status_file(&state_dir, &output_root),
)
.unwrap();
fs::write(daemon_pid_path(&state_dir), std::process::id().to_string()).unwrap();
let conflict = live_daemon_resource_conflict(&state_dir, &output_root)
.unwrap()
.unwrap();
assert_eq!(conflict.pid, Some(std::process::id()));
assert_eq!(
conflict.target_root,
normalized_abs_path(&output_root).unwrap()
);
}
#[test]
fn live_daemon_resource_conflict_allows_different_output() {
let temp = tempfile::TempDir::new().unwrap();
let state_dir = temp.path().join("state");
let daemon_output = temp.path().join("daemon-resources");
let manual_output = temp.path().join("manual-resources");
fs::create_dir_all(&daemon_output).unwrap();
fs::create_dir_all(&manual_output).unwrap();
write_daemon_status_file(
&daemon_status_path(&state_dir),
&test_daemon_status_file(&state_dir, &daemon_output),
)
.unwrap();
fs::write(daemon_pid_path(&state_dir), std::process::id().to_string()).unwrap();
assert!(live_daemon_resource_conflict(&state_dir, &manual_output)
.unwrap()
.is_none());
}
#[test]
fn direct_write_rejects_live_daemon_for_same_output() {
let temp = tempfile::TempDir::new().unwrap();
let state_dir = temp.path().join("state");
let output_root = temp.path().join("resources");
fs::create_dir_all(&output_root).unwrap();
write_daemon_status_file(
&daemon_status_path(&state_dir),
&test_daemon_status_file(&state_dir, &output_root),
)
.unwrap();
fs::write(daemon_pid_path(&state_dir), std::process::id().to_string()).unwrap();
let options = CliOptions {
state_dir,
config: OfficialUpdateConfig {
output_root,
..OfficialUpdateConfig::default()
},
..CliOptions::default()
};
let error = assert_no_live_daemon_output_conflict(&options, "repair").unwrap_err();
assert!(error.to_string().contains("同一资源目录"));
assert!(error.to_string().contains("locked"));
}
#[test]
fn clean_stable_removes_corrupt_pid_and_lock_files() {
let temp = tempfile::TempDir::new().unwrap();
let output_root = temp.path().join("resources");
let state_dir = temp.path().join("state");
fs::create_dir_all(&output_root).unwrap();
fs::create_dir_all(&state_dir).unwrap();
let resource_lock = output_root.join(".official-sync.lock");
let control_lock = daemon_control_lock_path(&state_dir);
let pid_path = daemon_pid_path(&state_dir);
fs::write(&resource_lock, "not-a-pid").unwrap();
fs::write(&control_lock, "not-a-pid").unwrap();
fs::write(&pid_path, "not-a-pid").unwrap();
let options = CliOptions {
command: CliCommand::CleanStable,
output_format: OutputFormat::Json,
state_dir,
config: OfficialUpdateConfig {
output_root,
..OfficialUpdateConfig::default()
},
..CliOptions::default()
};
run_clean_stable_command(&options).unwrap();
assert!(!resource_lock.exists());
assert!(!control_lock.exists());
assert!(!pid_path.exists());
}
#[test]
fn rpc_stop_wait_reports_already_stopped_without_signal() {
assert!(!wait_for_rpc_stop_or_terminate(0, Duration::from_millis(1)).unwrap());
}
#[test]
fn refresh_rpc_selection_only_for_default_daemon_shape() {
let options = parse(&["bat", "refresh"]).unwrap();
assert!(refresh_should_use_daemon_rpc(&options, "refresh"));
let options = parse(&["bat", "refresh", "--force"]).unwrap();
assert!(refresh_should_use_daemon_rpc(&options, "refresh"));
let options = parse(&["bat", "refresh", "--output", "/tmp/other"]).unwrap();
assert!(!refresh_should_use_daemon_rpc(&options, "refresh"));
let options = parse(&["bat", "refresh", "--server-info-file", "ProdNotice.json"]).unwrap();
assert!(!refresh_should_use_daemon_rpc(&options, "refresh"));
let options = parse(&["bat", "repair"]).unwrap();
assert!(!refresh_should_use_daemon_rpc(&options, "repair"));
}
#[test]
fn explicit_no_quiet_up_to_date_overrides_watch_default() {
let options = parse(&[
"bat",
"--watch",
"--no-quiet-up-to-date",
"--interval",
"1h",
])
.unwrap();
assert!(options.watch);
assert!(!options.quiet_up_to_date);
assert!(options.quiet_up_to_date_explicit);
}
#[test]
fn rejects_watch_dry_run_and_zero_interval() {
let error = parse(&["bat", "--watch", "--dry-run"]).unwrap_err();
assert!(error.to_string().contains("不能和 --dry-run"));
let error = parse(&["bat", "--daemon", "--watch"]).unwrap_err();
assert!(error.to_string().contains("不要同时传 --watch"));
let error = parse(&["bat", "status", "--watch"]).unwrap_err();
assert!(error.to_string().contains("status/stop"));
let error = parse(&["bat", "status", "--output", "/tmp/resources"]).unwrap_err();
assert!(error.to_string().contains("--state-dir"));
let error = parse(&["bat", "--watch", "--interval-seconds", "0"]).unwrap_err();
assert!(error.to_string().contains("必须大于 0"));
let error = parse(&["bat", "--watch", "--error-retry-seconds", "0"]).unwrap_err();
assert!(error.to_string().contains("失败重试间隔"));
}
#[test]
fn parses_duration_units() {
assert_eq!(parse_duration("1h").unwrap(), Duration::from_secs(3600));
assert_eq!(parse_duration("60m").unwrap(), Duration::from_secs(3600));
assert_eq!(parse_duration("90s").unwrap(), Duration::from_secs(90));
assert_eq!(parse_duration("250ms").unwrap(), Duration::from_millis(250));
assert_eq!(parse_duration("3600").unwrap(), Duration::from_secs(3600));
}
#[test]
fn formats_duration_args_for_child_process() {
assert_eq!(format_duration_arg(Duration::from_secs(3600)), "1h");
assert_eq!(format_duration_arg(Duration::from_secs(90 * 60)), "90m");
assert_eq!(format_duration_arg(Duration::from_secs(45)), "45s");
assert_eq!(format_duration_arg(Duration::from_millis(250)), "250ms");
}
#[test]
fn scheduled_forced_refresh_uses_beijing_local_times() {
let day = 20_000;
assert_eq!(
next_forced_refresh_at_or_after(beijing_time(day, 2, 30, 0)),
beijing_time(day, 3, 0, 0)
);
assert_eq!(
next_forced_refresh_at_or_after(beijing_time(day, 3, 0, 0)),
beijing_time(day, 3, 0, 0)
);
assert_eq!(
next_forced_refresh_after(beijing_time(day, 3, 0, 0)),
beijing_time(day, 16, 0, 0)
);
assert_eq!(
next_forced_refresh_after(beijing_time(day, 18, 0, 0)),
beijing_time(day + 1, 3, 0, 0)
);
assert_eq!(
duration_until(
next_forced_refresh_at_or_after(beijing_time(day, 17, 45, 0)),
beijing_time(day, 17, 45, 0)
),
Duration::from_secs(15 * 60)
);
}
#[test]
fn formats_progress_durations() {
assert_eq!(format_duration(Duration::from_millis(7)), "00:00.007");
assert_eq!(format_duration(Duration::from_millis(65_432)), "01:05.432");
}
#[test]
fn quiet_up_to_date_suppresses_only_clean_reports() {
assert!(!should_print_status(OfficialUpdateStatus::UpToDate, true));
assert!(should_print_status(OfficialUpdateStatus::Downloaded, true));
assert!(should_print_status(
OfficialUpdateStatus::WouldDownload,
true
));
assert!(should_print_status(OfficialUpdateStatus::UpToDate, false));
}
#[test]
fn status_string_is_stable() {
assert_eq!(OfficialUpdateStatus::UpToDate.as_str(), "up_to_date");
assert_eq!(
OfficialUpdateStatus::WouldDownload.as_str(),
"would_download"
);
assert_eq!(OfficialUpdateStatus::Downloaded.as_str(), "downloaded");
}
fn beijing_time(day: u64, hour: u64, minute: u64, second: u64) -> SystemTime {
system_time_from_beijing_local_seconds(
day.saturating_mul(SECONDS_PER_DAY)
.saturating_add(hour.saturating_mul(60 * 60))
.saturating_add(minute.saturating_mul(60))
.saturating_add(second),
)
}
fn test_daemon_status_file(state_dir: &Path, output_root: &Path) -> DaemonStatusFile {
DaemonStatusFile {
version: DAEMON_STATUS_VERSION,
pid: std::process::id(),
state: "sleeping".to_string(),
resource_output_root: output_root.to_path_buf(),
state_dir: state_dir.to_path_buf(),
log_path: daemon_log_path(state_dir),
started_unix_seconds: unix_seconds_now(),
updated_unix_seconds: unix_seconds_now(),
last_update_status: Some("up_to_date".to_string()),
last_error: None,
next_retry_seconds: Some(60),
pending_scheduled_force: false,
next_forced_refresh_unix_seconds: None,
command: vec!["bat".to_string(), "--daemon-child".to_string()],
}
}
}