feat(daemon): 任务历史文件态持久化与 .env 无参启动(issue #1)

任务持久化:
- 任务历史落 <state-dir>/bat-tasks.json(版本化、0600 原子写、不跟随
  symlink),生命周期转换时 write-through;seq 持久化避免 pid 复用撞 ID
- daemon 重启恢复历史,中断时仍 queued/running 的任务标记 failed,
  新增错误码 TASK_INTERRUPTED(BAT-ERR-700005)
- 损坏文件改名 .corrupt 留证后从空历史开始;无法识别的记录计数跳过
- core 新增 ErrorCode::ALL 公开码表与 from_id 反查(持久化错误码往返)

.env 配置:
- 首次启动在二进制所在目录释放 .env 模板(0600、create_new 防竞态),
  之后每次启动加载为进程环境变量(不覆盖已存在变量)
- 优先级:CLI > 进程环境变量 > .env > 内置默认;BAT_* 键映射到
  CliOptions 默认值,BAT_WATCH/BAT_DAEMON 仅对无子命令 Run 生效且
  CLI 显式模式/dry-run 时让位;Redis 键预留(未接入)
- 工具/代理"非默认"判断改按 env 应用后基线,status/stop/logs 在
  .env 存在时不误判;BAT_SKIP_ENV_FILE=1 整体禁用

验证:新增 10 个单测(env 解析/优先级/守卫回归/持久化往返/中断标记/
损坏恢复/未知类型跳过);真机 e2e:.env 释放加载、daemon 纯 .env 启动、
catalog.refresh 任务带类型化错误码落盘并跨 daemon 重启恢复(含日志);
fmt / clippy --workspace --all-targets -D warnings / test --workspace 全绿

Co-Authored-By: Claude Fable 5 <noreply@anthropic.com>
This commit is contained in:
2026-07-17 21:15:47 -07:00
co-authored by Claude Fable 5
parent ec40ed4660
commit 3edbe7ccee
5 changed files with 908 additions and 66 deletions
+810 -11
View File
@@ -92,6 +92,7 @@ fn main() {
}
fn run() -> anyhow::Result<i32> {
bootstrap_env_file();
let options = parse_args()?;
if matches!(
options.command,
@@ -200,6 +201,9 @@ struct CliOptions {
progress: bool,
banner: bool,
tail_lines: usize,
/// 环境变量(含 .env)应用后、命令行解析前的配置快照。
/// 工具/代理"是否命令行显式传入"的判断以它为基线。
env_baseline_config: OfficialUpdateConfig,
}
impl Default for CliOptions {
@@ -222,6 +226,7 @@ impl Default for CliOptions {
progress: true,
banner: true,
tail_lines: 200,
env_baseline_config: OfficialUpdateConfig::default(),
}
}
}
@@ -275,7 +280,9 @@ fn run_watch(options: CliOptions) -> anyhow::Result<()> {
let sync_lock = Arc::new(Mutex::new(()));
let (_task_worker, _task_context, _rpc_server) = if let Some(control) = daemon_control.as_ref()
{
let registry = TaskRegistry::new();
// 任务历史持久化在 state dir(此前已通过 validate_runtime_state_dir 校验)。
let (registry, restore_summary) = TaskRegistry::with_persistence(&daemon_state_dir);
logger.log_text("daemon", format!("任务历史:{restore_summary}"));
let (task_tx, task_rx) = mpsc::channel::<TaskJob>();
let worker = {
let registry = registry.clone();
@@ -746,6 +753,143 @@ struct TaskStore {
tasks: HashMap<String, TaskRecord>,
order: Vec<String>,
seq: u64,
/// 任务历史持久化文件路径;`None` 表示纯内存(测试等非 daemon 场景)。
persist_path: Option<PathBuf>,
}
/// daemon 任务历史持久化文件名(位于 state dir 内,`0600` 原子写)。
const TASKS_FILE_NAME: &str = "bat-tasks.json";
/// 任务历史文件结构版本。
const TASKS_FILE_VERSION: u32 = 1;
/// 任务历史文件的持久化形态(版本化;daemon 重启后恢复任务历史用)。
#[derive(Debug, Serialize, Deserialize)]
struct PersistedTaskFile {
version: u32,
/// 任务 ID 序号计数器;恢复它避免 pid 复用时新任务与历史任务撞 ID。
seq: u64,
tasks: Vec<PersistedTaskRecord>,
}
#[derive(Debug, Serialize, Deserialize)]
struct PersistedTaskRecord {
id: String,
kind: String,
status: String,
#[serde(default, skip_serializing_if = "Option::is_none")]
stage: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
message: Option<String>,
created_at: u64,
updated_at: u64,
#[serde(default, skip_serializing_if = "Option::is_none")]
started_at: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
finished_at: Option<u64>,
/// `ApiError` 的序列化形态(code/kind/domain/location/message/retryable)。
#[serde(default, skip_serializing_if = "Option::is_none")]
error: Option<serde_json::Value>,
#[serde(default, skip_serializing_if = "Option::is_none")]
result: Option<serde_json::Value>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
log: Vec<String>,
}
/// 把持久化的任务类型映射回静态字符串;未识别(如未来版本新增)返回 `None`。
fn task_kind_static(kind: &str) -> Option<&'static str> {
match kind {
RPC_METHOD_RESOURCE_SYNC => Some(RPC_METHOD_RESOURCE_SYNC),
RPC_METHOD_RESOURCE_VERIFY => Some(RPC_METHOD_RESOURCE_VERIFY),
RPC_METHOD_CATALOG_REFRESH => Some(RPC_METHOD_CATALOG_REFRESH),
_ => None,
}
}
/// 把持久化的任务状态映射回静态字符串;未识别返回 `None`。
fn task_status_static(status: &str) -> Option<&'static str> {
match status {
"queued" => Some("queued"),
"running" => Some("running"),
"succeeded" => Some("succeeded"),
"failed" => Some("failed"),
"cancelled" => Some("cancelled"),
_ => None,
}
}
impl PersistedTaskRecord {
fn from_record(record: &TaskRecord) -> Self {
Self {
id: record.id.clone(),
kind: record.kind.to_string(),
status: record.status.to_string(),
stage: record.stage.clone(),
message: record.message.clone(),
created_at: record.created_at,
updated_at: record.updated_at,
started_at: record.started_at,
finished_at: record.finished_at,
error: record
.error
.as_ref()
.and_then(|error| serde_json::to_value(error).ok()),
result: record.result.clone(),
log: record.log.clone(),
}
}
/// 还原为内存任务记录;kind/status 未识别时返回 `None`(调用方计数跳过)。
fn into_record(self) -> Option<TaskRecord> {
let kind = task_kind_static(&self.kind)?;
let status = task_status_static(&self.status)?;
// 错误从序列化形态还原:code 经码表反查(未登记回退 internal),
// location 固定为任务执行器(当前全部任务错误的唯一来源)。
let error = self.error.as_ref().map(|value| {
let code = value
.get("code")
.and_then(serde_json::Value::as_str)
.and_then(ErrorCode::from_id)
.unwrap_or(ErrorCode::INTERNAL);
let message = value
.get("message")
.and_then(serde_json::Value::as_str)
.unwrap_or("<持久化错误信息缺失>")
.to_string();
ApiError::new(code, "task.executor", message)
});
Some(TaskRecord {
id: self.id,
kind,
status,
stage: self.stage,
message: self.message,
created_at: self.created_at,
updated_at: self.updated_at,
started_at: self.started_at,
finished_at: self.finished_at,
error,
result: self.result,
cancel: Arc::new(AtomicBool::new(false)),
log: self.log,
})
}
}
/// 读取任务历史文件。文件缺失返回 `Ok(None)`;symlink、解析失败或版本不支持返回 `Err`。
fn load_persisted_tasks(path: &Path) -> Result<Option<PersistedTaskFile>, String> {
let Some(bytes) = read_file_no_symlink(path, "任务历史")? else {
return Ok(None);
};
let file: PersistedTaskFile = serde_json::from_slice(&bytes)
.map_err(|error| format!("解析任务历史失败 {}{error}", path.display()))?;
if file.version != TASKS_FILE_VERSION {
return Err(format!(
"不支持的任务历史版本 {},文件 {}",
file.version,
path.display()
));
}
Ok(Some(file))
}
/// 任务注册表句柄:包住内存存储,供 RPC handler 与 worker 共享。
@@ -757,16 +901,90 @@ struct TaskRegistry {
}
impl TaskRegistry {
/// 纯内存注册表(无持久化);生产 daemon 走 [`Self::with_persistence`]。
#[cfg(test)]
fn new() -> Self {
Self {
inner: Arc::new(Mutex::new(TaskStore {
tasks: HashMap::new(),
order: Vec::new(),
seq: 0,
persist_path: None,
})),
}
}
/// 从 state dir 恢复任务历史并启用持久化。
///
/// 中断时仍处于 queued/running 的任务标记为 `failed``TASK_INTERRUPTED`);
/// 文件缺失按空历史处理;文件损坏或版本不支持时改名 `.corrupt` 留证并从
/// 空历史开始。返回注册表与恢复摘要(供 daemon 日志记录)。
fn with_persistence(state_dir: &Path) -> (Self, String) {
let path = state_dir.join(TASKS_FILE_NAME);
let now = unix_seconds_now();
let mut seq = 0;
let mut tasks = HashMap::new();
let mut order = Vec::new();
let summary = match load_persisted_tasks(&path) {
Ok(None) => "无历史任务文件,从空任务历史开始".to_string(),
Ok(Some(file)) => {
seq = file.seq;
let total = file.tasks.len();
let mut interrupted = 0usize;
let mut skipped = 0usize;
for persisted in file.tasks {
let Some(mut record) = persisted.into_record() else {
skipped += 1;
continue;
};
if !record.is_finished() {
interrupted += 1;
record.status = "failed";
record.finished_at = Some(now);
record.updated_at = now;
record.error = Some(ApiError::new(
ErrorCode::TASK_INTERRUPTED,
"task.executor",
"daemon 停止/重启导致任务中断",
));
record
.log
.push("[daemon] 任务因 daemon 停止/重启而中断".to_string());
}
if tasks.insert(record.id.clone(), record.clone()).is_none() {
order.push(record.id);
} else {
skipped += 1;
}
}
format!("恢复任务历史 {total} 条(标记中断 {interrupted} 条,跳过无法识别 {skipped} 条)")
}
Err(error) => {
// 保留损坏文件供诊断(改名而非覆盖),从空历史开始。
let corrupt = path.with_extension("json.corrupt");
if fs::rename(&path, &corrupt).is_ok() {
format!(
"任务历史不可用({error});原文件已改名保留为 {}",
corrupt.display()
)
} else {
format!("任务历史不可用({error});且无法改名保留原文件")
}
}
};
let registry = Self {
inner: Arc::new(Mutex::new(TaskStore {
tasks,
order,
seq,
persist_path: Some(path),
})),
};
// 把中断标记(或空历史)立即写回,保证文件与内存视图一致。
registry.lock().persist();
(registry, summary)
}
fn lock(&self) -> std::sync::MutexGuard<'_, TaskStore> {
self.inner
.lock()
@@ -797,14 +1015,23 @@ impl TaskRegistry {
store.tasks.insert(id.clone(), record);
store.order.push(id.clone());
store.prune();
store.persist();
id
}
fn update<F: FnOnce(&mut TaskRecord)>(&self, id: &str, update: F) {
let mut store = self.lock();
let mut status_changed = false;
if let Some(record) = store.tasks.get_mut(id) {
let previous_status = record.status;
update(record);
record.updated_at = unix_seconds_now();
status_changed = record.status != previous_status;
}
// 只在生命周期转换时落盘;stage/message/log 的高频进度更新以内存为准,
// 随下一次转换一起写入(避免每个进度事件一次磁盘写)。
if status_changed {
store.persist();
}
}
@@ -863,6 +1090,35 @@ impl TaskRegistry {
}
impl TaskStore {
/// 把当前任务历史落盘(`0600` 原子写、不跟随 symlink)。
///
/// 持久化未启用时为 no-op;写失败只记 stderr(进 daemon 日志),
/// 不让持久化故障拖垮任务执行本身。
fn persist(&self) {
let Some(path) = &self.persist_path else {
return;
};
let file = PersistedTaskFile {
version: TASKS_FILE_VERSION,
seq: self.seq,
tasks: self
.order
.iter()
.filter_map(|id| self.tasks.get(id))
.map(PersistedTaskRecord::from_record)
.collect(),
};
match serde_json::to_vec_pretty(&file) {
Ok(bytes) => {
if let Err(error) = write_file_atomic(path, &bytes, PRIVATE_FILE_MODE, "任务历史")
{
eprintln!("[daemon] 任务历史落盘失败:{error}");
}
}
Err(error) => eprintln!("[daemon] 任务历史序列化失败:{error}"),
}
}
/// 裁剪最旧的已结束任务,把内存占用控制在上限内;运行中/排队中的任务不裁剪。
fn prune(&mut self) {
while self.order.len() > MAX_RETAINED_TASKS {
@@ -2630,7 +2886,7 @@ fn run_daemon_restart(options: &CliOptions, command_name: &'static str) -> anyho
let has_explicit_options = options.sync_option_explicit
|| options.output_explicit
|| options.proxy_option_explicit
|| tools_are_non_default(&options.config);
|| tools_are_non_default(&options.config, &options.env_baseline_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)?;
@@ -4673,14 +4929,275 @@ fn should_print_status(status: OfficialUpdateStatus, quiet_up_to_date: bool) ->
!(quiet_up_to_date && status == OfficialUpdateStatus::UpToDate)
}
fn parse_args() -> anyhow::Result<CliOptions> {
parse_args_from(env::args())
/// `.env` 配置文件名(位于 bat 二进制所在目录)。
const ENV_FILE_NAME: &str = ".env";
/// 设为 `1` 时完全跳过 `.env` 的生成与加载(测试与特殊部署场景用)。
const SKIP_ENV_FILE_VAR: &str = "BAT_SKIP_ENV_FILE";
/// 首次启动释放的 `.env` 配置模板。
const ENV_TEMPLATE: &str = r#"# BlueArchive Toolkit 配置文件(bat 首次启动自动生成)
#
# `bat`
# > > >
# 1/0/true/false/yes/no/on/off
# BAT_SKIP_ENV_FILE=1 bat
# ---- ----
# ./bat-resources
BAT_OUTPUT=./bat-resources
# app-version / connection-group / server-info 1
BAT_AUTO_DISCOVER=1
# bat.sock / / /tmp/bat-pid
#BAT_STATE_DIR=/tmp/bat-pid
# watch daemon
# `bat` 1 daemon
#BAT_WATCH=0
#BAT_DAEMON=0
#
#BAT_INTERVAL_SECONDS=3600
#BAT_ERROR_RETRY_SECONDS=60
# ---- ----
# URL http/https/socks4/socks4a/socks5/socks5h
# HTTPS_PROXY / ALL_PROXY / HTTP_PROXY
#BAT_PROXY=http://127.0.0.1:7897
# 1
#BAT_NO_PROXY=0
# ---- ----
#BAT_APP_VERSION=
#BAT_CONNECTION_GROUP=
#BAT_LAUNCHER_VERSION=
# windows,android
#BAT_PLATFORMS=windows,android
#BAT_CURL=curl
#BAT_UNZIP=unzip
# ---- ----
# 1 JSON
#BAT_JSON=0
# watch/daemon 1
#BAT_QUIET_UP_TO_DATE=
# ---- Redis----
# <BAT_STATE_DIR>/bat-tasks.json
# Redis
#BAT_REDIS_URL=redis://127.0.0.1:6379
#BAT_REDIS_PASSWORD=
"#;
/// `.env` 引导:首次启动时在二进制所在目录释放配置模板,之后每次启动把其中的
/// 键加载为进程环境变量(不覆盖已存在的环境变量,保持"环境变量 > .env"优先级)。
///
/// 任何失败只在 stderr 警告、不中断启动——`.env` 是便利层,不是启动硬依赖。
fn bootstrap_env_file() {
if env::var(SKIP_ENV_FILE_VAR).map(|value| value == "1") == Ok(true) {
return;
}
let Ok(exe_path) = env::current_exe() else {
return;
};
let Some(exe_dir) = exe_path.parent() else {
return;
};
let path = exe_dir.join(ENV_FILE_NAME);
if !path.exists() {
match write_env_template(&path) {
Ok(()) => eprintln!(
"已生成配置模板 {}(编辑其中的 BAT_* 配置后,直接运行 `bat` 即可按 .env 启动)",
path.display()
),
Err(error) => {
eprintln!("警告:生成 .env 配置模板失败 {}{error}", path.display());
return;
}
}
}
match fs::read_to_string(&path) {
Ok(content) => apply_env_file(&content),
Err(error) => eprintln!("警告:读取 .env 失败 {}{error}", path.display()),
}
}
/// 以 `create_new` 原子创建模板文件,避免并发启动时互相覆盖;unix 下限制 `0600`
/// 权限(`.env` 可能保存代理凭据等敏感配置)。
fn write_env_template(path: &Path) -> std::io::Result<()> {
let mut open_options = OpenOptions::new();
open_options.write(true).create_new(true);
#[cfg(unix)]
open_options.mode(PRIVATE_FILE_MODE);
let mut file = open_options.open(path)?;
file.write_all(ENV_TEMPLATE.as_bytes())
}
/// 解析 `.env` 内容,把进程环境里尚不存在的键设为环境变量。
fn apply_env_file(content: &str) {
for (line_number, raw_line) in content.lines().enumerate() {
let line = raw_line.trim();
if line.is_empty() || line.starts_with('#') {
continue;
}
let Some((key, value)) = parse_env_line(line) else {
eprintln!(
"警告:.env 第 {} 行无法解析,已忽略:{raw_line}",
line_number + 1
);
continue;
};
if env::var_os(&key).is_none() {
env::set_var(&key, value);
}
}
}
/// 解析单行 `KEY=VALUE`。key 须为 `[A-Za-z_][A-Za-z0-9_]*`;值两侧的成对
/// 单/双引号会剥除。不支持 `export` 前缀和多行值。
fn parse_env_line(line: &str) -> Option<(String, String)> {
let (key, value) = line.split_once('=')?;
let key = key.trim();
let valid_key = !key.is_empty()
&& key.chars().enumerate().all(|(index, character)| {
character == '_'
|| character.is_ascii_alphabetic()
|| (index > 0 && character.is_ascii_digit())
});
if !valid_key {
return None;
}
let mut value = value.trim();
if value.len() >= 2 {
let bytes = value.as_bytes();
let quoted = (bytes[0] == b'"' && bytes[value.len() - 1] == b'"')
|| (bytes[0] == b'\'' && bytes[value.len() - 1] == b'\'');
if quoted {
value = &value[1..value.len() - 1];
}
}
Some((key.to_string(), value.to_string()))
}
/// `.env`/环境变量提供的运行模式开关(延迟到命令确定后应用;值型配置
/// 由 [`apply_bat_env_overrides`] 直接写入 options)。
struct EnvModeOverrides {
watch: bool,
daemon: bool,
}
/// 把 `BAT_*` 环境变量作为配置默认值写入 options。
///
/// 不标记任何 `*_explicit`(命令行参数在其后解析、总是覆盖);非法值报错
/// 而非静默忽略,保证配置问题可诊断。
fn apply_bat_env_overrides(
options: &mut CliOptions,
env_lookup: &impl Fn(&str) -> Option<String>,
) -> anyhow::Result<EnvModeOverrides> {
fn parse_env_bool(key: &str, value: &str) -> anyhow::Result<bool> {
match value.to_ascii_lowercase().as_str() {
"1" | "true" | "yes" | "on" => Ok(true),
"0" | "false" | "no" | "off" => Ok(false),
other => Err(anyhow::anyhow!(
"环境变量 {key} 的布尔值无效:{other}(支持 1/0/true/false/yes/no/on/off"
)),
}
}
fn parse_env_seconds(key: &str, value: &str) -> anyhow::Result<Duration> {
let seconds = value
.parse::<u64>()
.map_err(|error| anyhow::anyhow!("环境变量 {key} 的秒数无效:{error}"))?;
Ok(Duration::from_secs(seconds))
}
// 空值视为未设置:模板里保留 `BAT_XXX=` 形式的空行不产生副作用。
let value = |key: &str| {
env_lookup(key)
.map(|value| value.trim().to_string())
.filter(|value| !value.is_empty())
};
if let Some(v) = value("BAT_OUTPUT") {
options.config.output_root = PathBuf::from(v);
}
if let Some(v) = value("BAT_STATE_DIR") {
options.state_dir = PathBuf::from(v);
}
if let Some(v) = value("BAT_AUTO_DISCOVER") {
options.config.auto_discover = parse_env_bool("BAT_AUTO_DISCOVER", &v)?;
}
if let Some(v) = value("BAT_APP_VERSION") {
options.config.app_version = Some(v);
}
if let Some(v) = value("BAT_CONNECTION_GROUP") {
options.config.connection_group = Some(v);
}
if let Some(v) = value("BAT_LAUNCHER_VERSION") {
options.config.launcher_version = v;
}
if let Some(v) = value("BAT_PLATFORMS") {
options.config.platforms = Some(parse_platforms(&v).map_err(anyhow::Error::msg)?);
}
if let Some(v) = value("BAT_CURL") {
options.config.curl_command = PathBuf::from(v);
}
if let Some(v) = value("BAT_UNZIP") {
options.config.unzip_command = PathBuf::from(v);
}
if let Some(v) = value("BAT_PROXY") {
options.config.curl_proxy = parse_proxy_config(&v)?;
}
if let Some(v) = value("BAT_NO_PROXY") {
if parse_env_bool("BAT_NO_PROXY", &v)? {
options.config.curl_proxy = CurlProxyConfig::disabled();
}
}
if let Some(v) = value("BAT_INTERVAL_SECONDS") {
options.interval = parse_env_seconds("BAT_INTERVAL_SECONDS", &v)?;
}
if let Some(v) = value("BAT_ERROR_RETRY_SECONDS") {
options.error_retry_interval = parse_env_seconds("BAT_ERROR_RETRY_SECONDS", &v)?;
}
if let Some(v) = value("BAT_JSON") {
if parse_env_bool("BAT_JSON", &v)? {
options.output_format = OutputFormat::Json;
}
}
if let Some(v) = value("BAT_QUIET_UP_TO_DATE") {
options.quiet_up_to_date = parse_env_bool("BAT_QUIET_UP_TO_DATE", &v)?;
options.quiet_up_to_date_explicit = true;
}
let watch = match value("BAT_WATCH") {
Some(v) => parse_env_bool("BAT_WATCH", &v)?,
None => false,
};
let daemon = match value("BAT_DAEMON") {
Some(v) => parse_env_bool("BAT_DAEMON", &v)?,
None => false,
};
Ok(EnvModeOverrides { watch, daemon })
}
fn parse_args() -> anyhow::Result<CliOptions> {
parse_args_with_env(env::args(), |key| env::var(key).ok())
}
/// 测试入口:不读环境变量,解析结果只由参数决定。
#[cfg(test)]
fn parse_args_from(raw_args: impl IntoIterator<Item = String>) -> anyhow::Result<CliOptions> {
parse_args_with_env(raw_args, |_| None)
}
fn parse_args_with_env(
raw_args: impl IntoIterator<Item = String>,
env_lookup: impl Fn(&str) -> Option<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();
// `BAT_*` 环境变量(含 .env 加载的)先作为默认值写入,不标记 explicit;
// 命令行参数随后解析,逐字段覆盖。工具/代理的"非默认"判断以本基线为准,
// 保证 status/stop/logs 在 .env 存在时不误判为显式传了同步参数。
let env_modes = apply_bat_env_overrides(&mut options, &env_lookup)?;
options.env_baseline_config = options.config.clone();
let mut mode_flag_from_cli = false;
while let Some(flag) = args.next() {
match flag.as_str() {
@@ -4832,15 +5349,18 @@ fn parse_args_from(raw_args: impl IntoIterator<Item = String>) -> anyhow::Result
"--watch" => {
options.watch = true;
options.sync_option_explicit = true;
mode_flag_from_cli = true;
}
"--daemon" => {
options.daemon = true;
options.sync_option_explicit = true;
mode_flag_from_cli = true;
}
"--daemon-child" => {
options.daemon_child = true;
options.watch = true;
options.sync_option_explicit = true;
mode_flag_from_cli = true;
}
"--interval" => {
options.interval = parse_duration(&next_option_value(&mut args, &flag)?)?;
@@ -4911,12 +5431,24 @@ fn parse_args_from(raw_args: impl IntoIterator<Item = String>) -> anyhow::Result
}
}
// `BAT_WATCH` / `BAT_DAEMON` 只影响无子命令的 Run(无参启动场景);
// 命令行显式选择了运行模式或 dry-run 时让位(命令行优先于 .env),
// status/verify 等子命令不受其影响。daemon 优先于 watchdaemon 自带 watch)。
if matches!(options.command, CliCommand::Run) && !mode_flag_from_cli && !options.config.dry_run
{
if env_modes.daemon {
options.daemon = true;
} else if env_modes.watch {
options.watch = true;
}
}
match options.command {
CliCommand::Status | CliCommand::Stop | CliCommand::Logs => {
if options.sync_option_explicit
|| options.output_explicit
|| options.proxy_option_explicit
|| tools_are_non_default(&options.config)
|| tools_are_non_default(&options.config, &options.env_baseline_config)
{
return Err(anyhow::anyhow!(
"status/stop/logs 只读取 --state-dir;资源同步参数和输出目录无关"
@@ -4953,7 +5485,7 @@ fn parse_args_from(raw_args: impl IntoIterator<Item = String>) -> anyhow::Result
if (options.sync_option_explicit
|| options.output_explicit
|| options.proxy_option_explicit
|| tools_are_non_default(&options.config))
|| tools_are_non_default(&options.config, &options.env_baseline_config))
&& !options.config.auto_discover
&& options.config.server_info_source.is_none()
&& options.config.connection_group.is_none()
@@ -5019,11 +5551,12 @@ fn ensure_command_not_set(command: CliCommand, next: &str) -> anyhow::Result<()>
}
}
fn tools_are_non_default(config: &OfficialUpdateConfig) -> bool {
let defaults = OfficialUpdateConfig::default();
config.curl_command != defaults.curl_command
|| config.curl_proxy != defaults.curl_proxy
|| config.unzip_command != defaults.unzip_command
/// 判断 curl/代理/unzip 是否偏离基线。基线是环境变量(含 .env)应用后的
/// 配置快照,因此只有命令行显式传入才算"非默认"。
fn tools_are_non_default(config: &OfficialUpdateConfig, baseline: &OfficialUpdateConfig) -> bool {
config.curl_command != baseline.curl_command
|| config.curl_proxy != baseline.curl_proxy
|| config.unzip_command != baseline.unzip_command
}
fn next_option_value(
@@ -5213,6 +5746,272 @@ mod tests {
parse_args_from(values.iter().map(|value| value.to_string()))
}
fn parse_with_env(values: &[&str], env: &[(&str, &str)]) -> anyhow::Result<CliOptions> {
let map: HashMap<String, String> = env
.iter()
.map(|(key, value)| (key.to_string(), value.to_string()))
.collect();
parse_args_with_env(values.iter().map(|value| value.to_string()), move |key| {
map.get(key).cloned()
})
}
#[test]
fn env_defaults_apply_and_cli_overrides() {
let options = parse_with_env(
&["bat"],
&[
("BAT_OUTPUT", "/srv/bat"),
("BAT_AUTO_DISCOVER", "1"),
("BAT_STATE_DIR", "/srv/state"),
("BAT_INTERVAL_SECONDS", "120"),
],
)
.unwrap();
assert_eq!(options.config.output_root, PathBuf::from("/srv/bat"));
assert!(options.config.auto_discover);
assert_eq!(options.state_dir, PathBuf::from("/srv/state"));
assert_eq!(options.interval, Duration::from_secs(120));
// 命令行覆盖环境变量。
let options = parse_with_env(
&["bat", "--output", "/cli/out"],
&[("BAT_OUTPUT", "/srv/bat")],
)
.unwrap();
assert_eq!(options.config.output_root, PathBuf::from("/cli/out"));
// 空值视为未设置(模板中保留 `BAT_XXX=` 空行无副作用)。
let options = parse_with_env(&["bat"], &[("BAT_OUTPUT", "")]).unwrap();
assert_eq!(
options.config.output_root,
CliOptions::default().config.output_root
);
}
#[test]
fn env_watch_daemon_only_affect_bare_run() {
let options = parse_with_env(&["bat"], &[("BAT_WATCH", "1")]).unwrap();
assert!(options.watch);
// 同时设置时 daemon 优先(daemon 自带 watch)。
let options = parse_with_env(&["bat"], &[("BAT_WATCH", "1"), ("BAT_DAEMON", "1")]).unwrap();
assert!(options.daemon);
assert!(!options.watch);
// 子命令不受影响:verify 内部置 dry-run,若误吃 env watch 会直接解析失败。
let options = parse_with_env(&["bat", "verify"], &[("BAT_WATCH", "1")]).unwrap();
assert!(!options.watch);
assert!(!options.daemon);
// 命令行显式 --dry-run 时 env watch 让位(命令行意图优先)。
let options = parse_with_env(&["bat", "--dry-run"], &[("BAT_WATCH", "1")]).unwrap();
assert!(!options.watch);
// 命令行显式选择前台 watch 时 env daemon 让位。
let options = parse_with_env(&["bat", "--watch"], &[("BAT_DAEMON", "1")]).unwrap();
assert!(options.watch);
assert!(!options.daemon);
}
#[test]
fn env_values_do_not_break_status_and_reload_guard() {
// .env 提供的代理/工具/输出目录不算"显式同步参数",status 应照常可用。
let options = parse_with_env(
&["bat", "status"],
&[
("BAT_PROXY", "http://127.0.0.1:7897"),
("BAT_CURL", "/usr/bin/curl"),
("BAT_OUTPUT", "/srv/bat"),
],
)
.unwrap();
assert_eq!(options.command, CliCommand::Status);
assert!(!tools_are_non_default(
&options.config,
&options.env_baseline_config
));
// 命令行再改工具才算偏离基线。
let options = parse_with_env(
&["bat", "--curl", "/opt/curl"],
&[("BAT_PROXY", "http://127.0.0.1:7897")],
)
.unwrap();
assert!(tools_are_non_default(
&options.config,
&options.env_baseline_config
));
}
#[test]
fn env_invalid_values_error() {
assert!(parse_with_env(&["bat"], &[("BAT_WATCH", "maybe")]).is_err());
assert!(parse_with_env(&["bat"], &[("BAT_INTERVAL_SECONDS", "abc")]).is_err());
assert!(parse_with_env(&["bat"], &[("BAT_PROXY", "ftp://x")]).is_err());
}
#[test]
fn parse_env_line_handles_quotes_and_rejects_bad_keys() {
assert_eq!(
parse_env_line("KEY=value"),
Some(("KEY".to_string(), "value".to_string()))
);
assert_eq!(
parse_env_line("KEY=\"quoted value\""),
Some(("KEY".to_string(), "quoted value".to_string()))
);
assert_eq!(
parse_env_line("KEY='single'"),
Some(("KEY".to_string(), "single".to_string()))
);
assert_eq!(
parse_env_line("BAT_OUTPUT = ./x"),
Some(("BAT_OUTPUT".to_string(), "./x".to_string()))
);
assert_eq!(parse_env_line("no_equals_sign"), None);
assert_eq!(parse_env_line("1BAD=x"), None);
assert_eq!(parse_env_line("BAD KEY=x"), None);
}
#[test]
fn env_template_is_parseable_and_bootstrap_ready() {
// 模板每个非注释行必须可解析;无参启动所需的最小配置默认启用。
let mut keys = Vec::new();
for line in ENV_TEMPLATE.lines() {
let line = line.trim();
if line.is_empty() || line.starts_with('#') {
continue;
}
let (key, _) =
parse_env_line(line).unwrap_or_else(|| panic!("模板行必须可解析:{line}"));
keys.push(key);
}
assert!(keys.contains(&"BAT_OUTPUT".to_string()));
assert!(keys.contains(&"BAT_AUTO_DISCOVER".to_string()));
}
#[test]
fn task_persistence_round_trip_and_interrupt_marking() {
let temp = tempfile::TempDir::new().unwrap();
let state_dir = temp.path();
let (registry, summary) = TaskRegistry::with_persistence(state_dir);
assert!(summary.contains("空任务历史"), "{summary}");
let finished_id = registry.create(TaskKind::Sync);
registry.update(&finished_id, |record| {
record.status = "running";
record.started_at = Some(record.created_at);
});
registry.append_log(&finished_id, "进度 1".to_string());
registry.update(&finished_id, |record| {
record.status = "succeeded";
record.finished_at = Some(record.created_at + 1);
record.result = Some(serde_json::json!({ "ok": true }));
});
let running_id = registry.create(TaskKind::Verify);
registry.update(&running_id, |record| record.status = "running");
let failed_id = registry.create(TaskKind::Refresh);
registry.update(&failed_id, |record| {
record.status = "failed";
record.error = Some(ApiError::new(
ErrorCode::HTTP_NOT_FOUND,
"task.executor",
"404",
));
});
drop(registry);
// 重启:恢复历史;running 任务标记中断;错误码经持久化往返保留。
let (registry, summary) = TaskRegistry::with_persistence(state_dir);
assert!(summary.contains("恢复任务历史 3 条"), "{summary}");
assert!(summary.contains("标记中断 1 条"), "{summary}");
let finished = registry.get(&finished_id).unwrap();
assert_eq!(finished.status, "succeeded");
assert_eq!(finished.result, Some(serde_json::json!({ "ok": true })));
assert_eq!(
registry.logs(&finished_id).unwrap(),
vec!["进度 1".to_string()]
);
let interrupted = registry.get(&running_id).unwrap();
assert_eq!(interrupted.status, "failed");
assert_eq!(
interrupted.error.as_ref().unwrap().code(),
ErrorCode::TASK_INTERRUPTED
);
assert!(interrupted.finished_at.is_some());
let failed = registry.get(&failed_id).unwrap();
assert_eq!(
failed.error.as_ref().unwrap().code(),
ErrorCode::HTTP_NOT_FOUND
);
// seq 持久化:重启后(本测试内 pid 相同)新任务不与历史撞 ID。
let new_id = registry.create(TaskKind::Sync);
assert!(
[&finished_id, &running_id, &failed_id]
.iter()
.all(|id| **id != new_id),
"新任务 ID {new_id} 与历史撞号"
);
assert_eq!(registry.list().len(), 4);
}
#[test]
fn task_persistence_recovers_from_corrupt_file() {
let temp = tempfile::TempDir::new().unwrap();
fs::write(temp.path().join(TASKS_FILE_NAME), b"not-json").unwrap();
let (registry, summary) = TaskRegistry::with_persistence(temp.path());
assert!(summary.contains("任务历史不可用"), "{summary}");
// 损坏文件改名留证,不静默覆盖。
assert!(temp.path().join("bat-tasks.json.corrupt").exists());
assert!(registry.list().is_empty());
// 恢复后可正常写入新历史。
registry.create(TaskKind::Sync);
let bytes = fs::read(temp.path().join(TASKS_FILE_NAME)).unwrap();
let file: PersistedTaskFile = serde_json::from_slice(&bytes).unwrap();
assert_eq!(file.version, TASKS_FILE_VERSION);
assert_eq!(file.tasks.len(), 1);
assert_eq!(file.tasks[0].status, "queued");
}
#[test]
fn task_persistence_skips_unknown_kinds() {
let temp = tempfile::TempDir::new().unwrap();
let file = serde_json::json!({
"version": TASKS_FILE_VERSION,
"seq": 9,
"tasks": [
{
"id": "task-1-1",
"kind": "patch.apply",
"status": "succeeded",
"created_at": 1,
"updated_at": 2
},
{
"id": "task-1-2",
"kind": "resource.sync",
"status": "succeeded",
"created_at": 3,
"updated_at": 4
}
]
});
fs::write(
temp.path().join(TASKS_FILE_NAME),
serde_json::to_vec(&file).unwrap(),
)
.unwrap();
let (registry, summary) = TaskRegistry::with_persistence(temp.path());
assert!(summary.contains("跳过无法识别 1 条"), "{summary}");
assert!(registry.get("task-1-1").is_none());
assert_eq!(registry.get("task-1-2").unwrap().status, "succeeded");
}
#[test]
fn parses_auto_discover_sync_args() {
let options = parse(&[