ded3253b08
按四不写原则移除 PR #159 引入、已被策划 V2 取代的整条 V1 链路 删除 runtime_protocol 下六个 V1 模块(planning_storage/submit/approval/coordinator/hydrate/provider_usage) V1 与 V2 共用的 GDD 数据模型抽到新模块 planning_gdd_model.rs 供 V2 继续复用 删除 prompt manifest 中 planning Agent 目录、role overlay、plan sections 与 supervisorPlan 组合及对应生成常量 删除四个 V1 提示词文件(roles/project-planning.md、plan/common.md、plan/supervisor-identity.md、plan/supervisor-playbook.md) 删除 game-creator.config.json 与配置代码中的 planning 能力开关 删除 CLI --swarm-chat 的 --plan 入口与 swarm_cli 中的 plan source 分支 删除 provider_retry 的 planning session binding 与 plan 专用请求指纹链路 删除 provider_action_batch 的 plan.submit_gdd v4 批次形状校验与 planning 绑定字段 删除 tool_policy_snapshot / agent_native_tools / tool_plan_protocol 中的 plan 根阶段收窄与 plan.submit_gdd 身份门 删除 main_loop 的 plan_gdd blocker 投影、plan 信封修复回路与 plan submit 业务拒绝限流 删除 pending_recovery 与 recovery_scan 的 plan submit 锚点恢复、planning session 投影恢复与审批投影恢复 删除 acceptance_graph 的 Fast GDD 取证覆盖校验与 plan 根验收前置门 删除 agent_db 的 plan.provider_usage、plan.gdd_decided、plan_submit_gdd.committed 三条专用持久车道及其预留配额 删除 AgentRuntimeState 的 plan_submit_gdd_rejection_count 字段 runtime_tools 的委派、run_status、goal_contract 恢复为通用路径(移除 plan 根对称性守卫与 acceptance gate 钩子) 同步删除只覆盖 V1 行为的测试用例(planning 澄清、锚点恢复、plan 根委派门、plan 提示词组合等)
818 lines
30 KiB
Rust
818 lines
30 KiB
Rust
use std::collections::BTreeSet;
|
|
use std::fs;
|
|
use std::path::{Path, PathBuf};
|
|
use std::time::{SystemTime, UNIX_EPOCH};
|
|
|
|
use serde::{Deserialize, Serialize};
|
|
use sha2::{Digest, Sha256};
|
|
|
|
use crate::agent::{
|
|
agent_runtime_json_sidecar_backup_path, read_agent_runtime_json_sidecar_with_max_bytes,
|
|
remove_agent_runtime_json_sidecar_backup, write_agent_runtime_json_sidecar_with_max_bytes,
|
|
};
|
|
|
|
pub(crate) const PROVIDER_RETRY_SCHEMA_VERSION: &str = "game-creator-provider-retry.v1";
|
|
|
|
const PROVIDER_RETRY_RELATIVE_DIRECTORY: &str = ".agent/runtime/provider-retries";
|
|
const PROVIDER_RETRY_SIDECAR_MAX_BYTES: usize = 32 * 1024;
|
|
const PROVIDER_RETRY_LABEL: &str = "Agent Runtime Provider 重试记录";
|
|
const SAFE_PATH_COMPONENT_MAX_CHARS: usize = 160;
|
|
const HASHED_PATH_PREFIX_MAX_CHARS: usize = 80;
|
|
|
|
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
|
|
#[serde(deny_unknown_fields, rename_all = "camelCase")]
|
|
pub(crate) struct AgentRuntimeProviderRetryIdentity {
|
|
pub(crate) project_id: String,
|
|
pub(crate) agent_id: String,
|
|
pub(crate) task_id: String,
|
|
pub(crate) session_id: String,
|
|
pub(crate) run_id: String,
|
|
pub(crate) source: String,
|
|
pub(crate) goal_id: Option<String>,
|
|
pub(crate) goal_revision: u64,
|
|
pub(crate) goal_snapshot_fingerprint: String,
|
|
pub(crate) applied_steer_cursor: u64,
|
|
pub(crate) request_kind: String,
|
|
pub(crate) base_request_slot: String,
|
|
pub(crate) request_fingerprint: String,
|
|
pub(crate) provider_config_fingerprint: String,
|
|
pub(crate) web_search_enabled: bool,
|
|
pub(crate) allow_idle_context_compaction: bool,
|
|
}
|
|
|
|
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
|
|
#[serde(deny_unknown_fields, rename_all = "camelCase")]
|
|
pub(crate) struct AgentRuntimeProviderRetryRecord {
|
|
pub(crate) schema_version: String,
|
|
pub(crate) identity: AgentRuntimeProviderRetryIdentity,
|
|
pub(crate) next_request_slot: String,
|
|
pub(crate) next_attempt: u32,
|
|
pub(crate) max_retries: u32,
|
|
pub(crate) backoff_ms: u64,
|
|
pub(crate) retry_at_ms: u64,
|
|
pub(crate) error_kind: String,
|
|
pub(crate) error_fingerprint: String,
|
|
pub(crate) created_at_ms: u64,
|
|
pub(crate) updated_at_ms: u64,
|
|
}
|
|
|
|
pub(crate) fn now_ms() -> u64 {
|
|
SystemTime::now()
|
|
.duration_since(UNIX_EPOCH)
|
|
.unwrap_or_default()
|
|
.as_millis()
|
|
.try_into()
|
|
.unwrap_or(u64::MAX)
|
|
}
|
|
|
|
pub(crate) fn remaining_ms(record: &AgentRuntimeProviderRetryRecord) -> u64 {
|
|
remaining_ms_at(record.retry_at_ms, now_ms())
|
|
}
|
|
|
|
pub(crate) fn read_for_run_at(
|
|
root: &Path,
|
|
agent_id: &str,
|
|
run_id: &str,
|
|
) -> Result<Option<AgentRuntimeProviderRetryRecord>, String> {
|
|
validate_path_identity(agent_id, run_id)?;
|
|
let relative_path = provider_retry_relative_path(agent_id, run_id);
|
|
let Some(record) = read_agent_runtime_json_sidecar_with_max_bytes(
|
|
root,
|
|
&relative_path,
|
|
PROVIDER_RETRY_LABEL,
|
|
PROVIDER_RETRY_SIDECAR_MAX_BYTES,
|
|
)?
|
|
else {
|
|
return Ok(None);
|
|
};
|
|
validate_record(&record)?;
|
|
if record.identity.agent_id != agent_id || record.identity.run_id != run_id {
|
|
return Err("Provider 重试记录与路径 Agent/run 身份冲突".to_string());
|
|
}
|
|
Ok(Some(record))
|
|
}
|
|
|
|
pub(crate) fn read_matching_at(
|
|
root: &Path,
|
|
identity: &AgentRuntimeProviderRetryIdentity,
|
|
) -> Result<Option<AgentRuntimeProviderRetryRecord>, String> {
|
|
validate_identity(identity)?;
|
|
let Some(record) = read_for_run_at(root, &identity.agent_id, &identity.run_id)? else {
|
|
return Ok(None);
|
|
};
|
|
if record.identity != *identity {
|
|
return Err("Provider 重试记录身份冲突".to_string());
|
|
}
|
|
Ok(Some(record))
|
|
}
|
|
|
|
pub(crate) fn list_at(root: &Path) -> Result<Vec<AgentRuntimeProviderRetryRecord>, String> {
|
|
let directory = root.join(PROVIDER_RETRY_RELATIVE_DIRECTORY);
|
|
let agent_entries = match fs::read_dir(&directory) {
|
|
Ok(entries) => entries,
|
|
Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(Vec::new()),
|
|
Err(error) => {
|
|
return Err(format!(
|
|
"读取 Provider 重试目录失败:{}: {error}",
|
|
directory.display()
|
|
));
|
|
}
|
|
};
|
|
let mut primary_paths = BTreeSet::new();
|
|
for (agent_index, agent_entry) in agent_entries.enumerate() {
|
|
if agent_index >= 1024 {
|
|
return Err("Provider 重试目录超过 1024 个 Agent 上限".to_string());
|
|
}
|
|
let agent_entry = agent_entry.map_err(|error| {
|
|
format!(
|
|
"读取 Provider 重试 Agent 目录项失败:{}: {error}",
|
|
directory.display()
|
|
)
|
|
})?;
|
|
let agent_type = agent_entry.file_type().map_err(|error| {
|
|
format!(
|
|
"读取 Provider 重试 Agent 目录类型失败:{}: {error}",
|
|
agent_entry.path().display()
|
|
)
|
|
})?;
|
|
if agent_type.is_symlink() || !agent_type.is_dir() {
|
|
return Err(format!(
|
|
"Provider 重试 Agent 项必须是普通目录:{}",
|
|
agent_entry.path().display()
|
|
));
|
|
}
|
|
let run_entries = fs::read_dir(agent_entry.path()).map_err(|error| {
|
|
format!(
|
|
"读取 Provider 重试 run 目录失败:{}: {error}",
|
|
agent_entry.path().display()
|
|
)
|
|
})?;
|
|
for (run_index, run_entry) in run_entries.enumerate() {
|
|
if run_index >= 1024 || primary_paths.len() >= 1024 {
|
|
return Err("Provider 重试记录超过 1024 条上限".to_string());
|
|
}
|
|
let run_entry = run_entry.map_err(|error| {
|
|
format!(
|
|
"读取 Provider 重试 run 目录项失败:{}: {error}",
|
|
agent_entry.path().display()
|
|
)
|
|
})?;
|
|
let run_type = run_entry.file_type().map_err(|error| {
|
|
format!(
|
|
"读取 Provider 重试 run 项类型失败:{}: {error}",
|
|
run_entry.path().display()
|
|
)
|
|
})?;
|
|
if run_type.is_symlink() || !run_type.is_file() {
|
|
return Err(format!(
|
|
"Provider 重试 run 项必须是普通文件:{}",
|
|
run_entry.path().display()
|
|
));
|
|
}
|
|
let file_name = run_entry
|
|
.file_name()
|
|
.into_string()
|
|
.map_err(|_| "Provider 重试文件名必须是 UTF-8".to_string())?;
|
|
let primary_path = if file_name.ends_with(".json") {
|
|
run_entry.path()
|
|
} else if let Some(primary_name) = file_name
|
|
.strip_prefix('.')
|
|
.and_then(|value| value.strip_suffix(".previous"))
|
|
.filter(|value| value.ends_with(".json"))
|
|
{
|
|
run_entry.path().with_file_name(primary_name)
|
|
} else {
|
|
return Err(format!("Provider 重试目录包含未知文件:{file_name}"));
|
|
};
|
|
primary_paths.insert(primary_path);
|
|
}
|
|
}
|
|
|
|
let mut records = Vec::with_capacity(primary_paths.len());
|
|
for primary_path in primary_paths {
|
|
let relative_path = primary_path.strip_prefix(root).map_err(|_| {
|
|
format!(
|
|
"Provider 重试文件不在项目目录内:{}",
|
|
primary_path.display()
|
|
)
|
|
})?;
|
|
let relative_path_text = relative_path
|
|
.iter()
|
|
.map(|component| {
|
|
component
|
|
.to_str()
|
|
.ok_or_else(|| "Provider 重试相对路径必须是 UTF-8".to_string())
|
|
})
|
|
.collect::<Result<Vec<_>, _>>()?
|
|
.join("/");
|
|
let record = read_agent_runtime_json_sidecar_with_max_bytes(
|
|
root,
|
|
&relative_path_text,
|
|
PROVIDER_RETRY_LABEL,
|
|
PROVIDER_RETRY_SIDECAR_MAX_BYTES,
|
|
)?
|
|
.ok_or_else(|| "Provider 重试目录项在扫描期间消失".to_string())?;
|
|
validate_record(&record)?;
|
|
let expected = PathBuf::from(provider_retry_relative_path(
|
|
&record.identity.agent_id,
|
|
&record.identity.run_id,
|
|
));
|
|
if relative_path != expected {
|
|
return Err(format!(
|
|
"Provider 重试记录身份与路径冲突:{}",
|
|
relative_path.display()
|
|
));
|
|
}
|
|
records.push(record);
|
|
}
|
|
records.sort_by(|left, right| {
|
|
left.identity
|
|
.agent_id
|
|
.cmp(&right.identity.agent_id)
|
|
.then_with(|| left.identity.run_id.cmp(&right.identity.run_id))
|
|
});
|
|
Ok(records)
|
|
}
|
|
|
|
#[allow(clippy::too_many_arguments)]
|
|
pub(crate) fn write_next_at(
|
|
root: &Path,
|
|
identity: &AgentRuntimeProviderRetryIdentity,
|
|
next_request_slot: &str,
|
|
next_attempt: u32,
|
|
max_retries: u32,
|
|
backoff_ms: u64,
|
|
error_kind: &str,
|
|
error_fingerprint: &str,
|
|
) -> Result<AgentRuntimeProviderRetryRecord, String> {
|
|
validate_identity(identity)?;
|
|
let existing = read_for_run_at(root, &identity.agent_id, &identity.run_id)?;
|
|
let created_at_ms;
|
|
let minimum_updated_at_ms;
|
|
if let Some(existing) = existing.as_ref() {
|
|
if existing.identity != *identity {
|
|
return Err("Provider 重试记录身份冲突".to_string());
|
|
}
|
|
if existing.next_attempt == next_attempt {
|
|
if retry_payload_matches(
|
|
existing,
|
|
next_request_slot,
|
|
max_retries,
|
|
backoff_ms,
|
|
error_kind,
|
|
error_fingerprint,
|
|
) {
|
|
return Ok(existing.clone());
|
|
}
|
|
return Err("Provider 重试记录同一 attempt 的内容冲突".to_string());
|
|
}
|
|
let expected_attempt = existing
|
|
.next_attempt
|
|
.checked_add(1)
|
|
.ok_or_else(|| "Provider 重试 attempt 溢出".to_string())?;
|
|
if next_attempt != expected_attempt || max_retries != existing.max_retries {
|
|
return Err("Provider 重试记录 attempt 或 maxRetries 冲突".to_string());
|
|
}
|
|
created_at_ms = existing.created_at_ms;
|
|
minimum_updated_at_ms = existing.updated_at_ms;
|
|
} else {
|
|
if next_attempt != 1 {
|
|
return Err("首条 Provider 重试记录的 nextAttempt 必须为 1".to_string());
|
|
}
|
|
created_at_ms = now_ms();
|
|
minimum_updated_at_ms = created_at_ms;
|
|
}
|
|
|
|
let updated_at_ms = now_ms().max(created_at_ms).max(minimum_updated_at_ms);
|
|
let retry_at_ms = updated_at_ms
|
|
.checked_add(backoff_ms)
|
|
.ok_or_else(|| "Provider 重试到期时间溢出".to_string())?;
|
|
let record = AgentRuntimeProviderRetryRecord {
|
|
schema_version: PROVIDER_RETRY_SCHEMA_VERSION.to_string(),
|
|
identity: identity.clone(),
|
|
next_request_slot: next_request_slot.to_string(),
|
|
next_attempt,
|
|
max_retries,
|
|
backoff_ms,
|
|
retry_at_ms,
|
|
error_kind: error_kind.to_string(),
|
|
error_fingerprint: error_fingerprint.to_string(),
|
|
created_at_ms,
|
|
updated_at_ms,
|
|
};
|
|
validate_record(&record)?;
|
|
let relative_path = provider_retry_relative_path(&identity.agent_id, &identity.run_id);
|
|
write_agent_runtime_json_sidecar_with_max_bytes(
|
|
root,
|
|
&relative_path,
|
|
PROVIDER_RETRY_LABEL,
|
|
&record,
|
|
PROVIDER_RETRY_SIDECAR_MAX_BYTES,
|
|
)?;
|
|
let persisted = read_matching_at(root, identity)?
|
|
.ok_or_else(|| "Provider 重试记录写入后不存在".to_string())?;
|
|
if persisted != record {
|
|
return Err("Provider 重试记录并发写入后内容冲突".to_string());
|
|
}
|
|
Ok(persisted)
|
|
}
|
|
|
|
pub(crate) fn remove_at(root: &Path, agent_id: &str, run_id: &str) -> Result<(), String> {
|
|
validate_path_identity(agent_id, run_id)?;
|
|
let path = provider_retry_path(root, agent_id, run_id);
|
|
let backup_path = agent_runtime_json_sidecar_backup_path(&path);
|
|
remove_agent_runtime_json_sidecar_backup(&backup_path, PROVIDER_RETRY_LABEL)?;
|
|
remove_agent_runtime_json_sidecar_backup(&path, PROVIDER_RETRY_LABEL)
|
|
}
|
|
|
|
#[cfg(test)]
|
|
pub(crate) fn force_provider_retry_due_for_test_at(
|
|
root: &Path,
|
|
identity: &AgentRuntimeProviderRetryIdentity,
|
|
) -> Result<AgentRuntimeProviderRetryRecord, String> {
|
|
let mut record = read_matching_at(root, identity)?
|
|
.ok_or_else(|| "待强制到期的 Provider 重试记录不存在".to_string())?;
|
|
record.retry_at_ms = record.created_at_ms;
|
|
record.updated_at_ms = now_ms().max(record.updated_at_ms);
|
|
validate_record(&record)?;
|
|
let relative_path = provider_retry_relative_path(&identity.agent_id, &identity.run_id);
|
|
write_agent_runtime_json_sidecar_with_max_bytes(
|
|
root,
|
|
&relative_path,
|
|
PROVIDER_RETRY_LABEL,
|
|
&record,
|
|
PROVIDER_RETRY_SIDECAR_MAX_BYTES,
|
|
)?;
|
|
let persisted = read_matching_at(root, identity)?
|
|
.ok_or_else(|| "Provider 重试记录强制到期后不存在".to_string())?;
|
|
if persisted != record {
|
|
return Err("Provider 重试记录强制到期后内容冲突".to_string());
|
|
}
|
|
Ok(persisted)
|
|
}
|
|
|
|
fn validate_path_identity(agent_id: &str, run_id: &str) -> Result<(), String> {
|
|
if agent_id.trim().is_empty() || run_id.trim().is_empty() {
|
|
return Err("Provider 重试路径的 Agent/run 身份不能为空".to_string());
|
|
}
|
|
Ok(())
|
|
}
|
|
|
|
pub(crate) fn validate_identity(
|
|
identity: &AgentRuntimeProviderRetryIdentity,
|
|
) -> Result<(), String> {
|
|
validate_path_identity(&identity.agent_id, &identity.run_id)?;
|
|
for (label, value) in [
|
|
("projectId", identity.project_id.as_str()),
|
|
("taskId", identity.task_id.as_str()),
|
|
("sessionId", identity.session_id.as_str()),
|
|
("source", identity.source.as_str()),
|
|
("requestKind", identity.request_kind.as_str()),
|
|
("baseRequestSlot", identity.base_request_slot.as_str()),
|
|
] {
|
|
if value.trim().is_empty() {
|
|
return Err(format!("Provider 重试身份 {label} 不能为空"));
|
|
}
|
|
}
|
|
if !is_sha256(&identity.request_fingerprint) {
|
|
return Err("Provider 重试身份的 requestFingerprint 无效".to_string());
|
|
}
|
|
if !is_sha256(&identity.provider_config_fingerprint) {
|
|
return Err("Provider 重试身份的 providerConfigFingerprint 无效".to_string());
|
|
}
|
|
match identity.goal_id.as_deref() {
|
|
Some(goal_id)
|
|
if !goal_id.trim().is_empty()
|
|
&& identity.goal_revision > 0
|
|
&& is_sha256(&identity.goal_snapshot_fingerprint) => {}
|
|
None if identity.goal_revision == 0 && identity.goal_snapshot_fingerprint.is_empty() => {}
|
|
_ => return Err("Provider 重试身份的 Goal 绑定无效".to_string()),
|
|
}
|
|
Ok(())
|
|
}
|
|
|
|
fn validate_record(record: &AgentRuntimeProviderRetryRecord) -> Result<(), String> {
|
|
if record.schema_version != PROVIDER_RETRY_SCHEMA_VERSION {
|
|
return Err(format!(
|
|
"不支持的 Provider 重试记录版本:{}",
|
|
record.schema_version
|
|
));
|
|
}
|
|
validate_identity(&record.identity)?;
|
|
if record.next_attempt == 0
|
|
|| record.max_retries == 0
|
|
|| record.next_attempt > record.max_retries
|
|
{
|
|
return Err("Provider 重试记录的 attempt/maxRetries 无效".to_string());
|
|
}
|
|
let expected_slot = format!(
|
|
"{}-transient-{}",
|
|
record.identity.base_request_slot, record.next_attempt
|
|
);
|
|
if record.next_request_slot != expected_slot {
|
|
return Err("Provider 重试记录的 nextRequestSlot 与身份不匹配".to_string());
|
|
}
|
|
if record.backoff_ms == 0
|
|
|| record.created_at_ms == 0
|
|
|| record.updated_at_ms < record.created_at_ms
|
|
|| record.retry_at_ms < record.created_at_ms
|
|
{
|
|
return Err("Provider 重试记录的退避或时间字段无效".to_string());
|
|
}
|
|
if record.error_kind.trim().is_empty() || !is_sha256(&record.error_fingerprint) {
|
|
return Err("Provider 重试记录的错误身份无效".to_string());
|
|
}
|
|
Ok(())
|
|
}
|
|
|
|
fn retry_payload_matches(
|
|
record: &AgentRuntimeProviderRetryRecord,
|
|
next_request_slot: &str,
|
|
max_retries: u32,
|
|
backoff_ms: u64,
|
|
error_kind: &str,
|
|
error_fingerprint: &str,
|
|
) -> bool {
|
|
record.next_request_slot == next_request_slot
|
|
&& record.max_retries == max_retries
|
|
&& record.backoff_ms == backoff_ms
|
|
&& record.error_kind == error_kind
|
|
&& record.error_fingerprint == error_fingerprint
|
|
}
|
|
|
|
fn is_sha256(value: &str) -> bool {
|
|
value.len() == 64 && value.bytes().all(|byte| byte.is_ascii_hexdigit())
|
|
}
|
|
|
|
fn remaining_ms_at(retry_at_ms: u64, at_ms: u64) -> u64 {
|
|
retry_at_ms.saturating_sub(at_ms)
|
|
}
|
|
|
|
fn provider_retry_relative_path(agent_id: &str, run_id: &str) -> String {
|
|
format!(
|
|
"{PROVIDER_RETRY_RELATIVE_DIRECTORY}/{}/{}.json",
|
|
provider_retry_path_component(agent_id, "agent"),
|
|
provider_retry_path_component(run_id, "run")
|
|
)
|
|
}
|
|
|
|
fn provider_retry_path(root: &Path, agent_id: &str, run_id: &str) -> PathBuf {
|
|
root.join(provider_retry_relative_path(agent_id, run_id))
|
|
}
|
|
|
|
fn provider_retry_path_component(value: &str, fallback: &str) -> String {
|
|
if is_safe_path_component(value) {
|
|
return value.to_string();
|
|
}
|
|
let readable = readable_path_prefix(value, fallback);
|
|
format!("{readable}--{:x}", Sha256::digest(value.as_bytes()))
|
|
}
|
|
|
|
fn is_safe_path_component(value: &str) -> bool {
|
|
if value.is_empty()
|
|
|| value.trim() != value
|
|
|| value == "."
|
|
|| value == ".."
|
|
|| value.chars().count() > SAFE_PATH_COMPONENT_MAX_CHARS
|
|
|| value.ends_with('.')
|
|
|| !value.chars().all(|character| {
|
|
character.is_ascii_alphanumeric()
|
|
|| character == '-'
|
|
|| character == '_'
|
|
|| character == '.'
|
|
})
|
|
{
|
|
return false;
|
|
}
|
|
!is_windows_reserved_component(value)
|
|
}
|
|
|
|
fn readable_path_prefix(value: &str, fallback: &str) -> String {
|
|
let normalized = value
|
|
.trim()
|
|
.chars()
|
|
.map(|character| {
|
|
if character.is_ascii_alphanumeric()
|
|
|| character == '-'
|
|
|| character == '_'
|
|
|| character == '.'
|
|
{
|
|
character
|
|
} else {
|
|
'-'
|
|
}
|
|
})
|
|
.take(HASHED_PATH_PREFIX_MAX_CHARS)
|
|
.collect::<String>();
|
|
let normalized = normalized.trim_matches(|character| character == '-' || character == '.');
|
|
if normalized.is_empty() || is_windows_reserved_component(normalized) {
|
|
fallback.to_string()
|
|
} else {
|
|
normalized.to_string()
|
|
}
|
|
}
|
|
|
|
fn is_windows_reserved_component(value: &str) -> bool {
|
|
let stem = value
|
|
.split('.')
|
|
.next()
|
|
.unwrap_or_default()
|
|
.to_ascii_uppercase();
|
|
matches!(stem.as_str(), "CON" | "PRN" | "AUX" | "NUL")
|
|
|| (stem.len() == 4
|
|
&& matches!(&stem[..3], "COM" | "LPT")
|
|
&& matches!(stem.as_bytes()[3], b'1'..=b'9'))
|
|
}
|
|
|
|
#[cfg(test)]
|
|
mod tests {
|
|
use std::fs;
|
|
|
|
use tempfile::tempdir;
|
|
|
|
use super::*;
|
|
|
|
fn identity(agent_id: &str, run_id: &str) -> AgentRuntimeProviderRetryIdentity {
|
|
AgentRuntimeProviderRetryIdentity {
|
|
project_id: "project-provider-retry".to_string(),
|
|
agent_id: agent_id.to_string(),
|
|
task_id: "task-provider-retry".to_string(),
|
|
session_id: "session-provider-retry".to_string(),
|
|
run_id: run_id.to_string(),
|
|
source: "agent-chat".to_string(),
|
|
goal_id: Some("goal-provider-retry".to_string()),
|
|
goal_revision: 3,
|
|
goal_snapshot_fingerprint: "a".repeat(64),
|
|
applied_steer_cursor: 2,
|
|
request_kind: "tool-plan".to_string(),
|
|
base_request_slot: "loop-2-repair-0".to_string(),
|
|
request_fingerprint: "d".repeat(64),
|
|
provider_config_fingerprint: "e".repeat(64),
|
|
web_search_enabled: true,
|
|
allow_idle_context_compaction: false,
|
|
}
|
|
}
|
|
|
|
fn write_first(
|
|
root: &Path,
|
|
identity: &AgentRuntimeProviderRetryIdentity,
|
|
) -> AgentRuntimeProviderRetryRecord {
|
|
write_next_at(
|
|
root,
|
|
identity,
|
|
"loop-2-repair-0-transient-1",
|
|
1,
|
|
3,
|
|
250,
|
|
"transport",
|
|
&"b".repeat(64),
|
|
)
|
|
.expect("write first Provider retry")
|
|
}
|
|
|
|
#[test]
|
|
fn provider_retry_round_trips_and_same_attempt_is_idempotent() {
|
|
let directory = tempdir().expect("create temp directory");
|
|
let identity = identity("design-director", "run-1");
|
|
let first = write_first(directory.path(), &identity);
|
|
let read = read_matching_at(directory.path(), &identity)
|
|
.expect("read matching Provider retry")
|
|
.expect("Provider retry exists");
|
|
assert_eq!(read, first);
|
|
assert_eq!(
|
|
list_at(directory.path()).expect("list Provider retries"),
|
|
vec![first.clone()]
|
|
);
|
|
let repeated = write_first(directory.path(), &identity);
|
|
assert_eq!(repeated, first);
|
|
|
|
let second = write_next_at(
|
|
directory.path(),
|
|
&identity,
|
|
"loop-2-repair-0-transient-2",
|
|
2,
|
|
3,
|
|
500,
|
|
"timeout",
|
|
&"c".repeat(64),
|
|
)
|
|
.expect("advance Provider retry");
|
|
assert_eq!(second.next_attempt, 2);
|
|
assert_eq!(second.created_at_ms, first.created_at_ms);
|
|
assert!(second.updated_at_ms >= first.updated_at_ms);
|
|
|
|
let serialized = fs::read_to_string(provider_retry_path(
|
|
directory.path(),
|
|
&identity.agent_id,
|
|
&identity.run_id,
|
|
))
|
|
.expect("read serialized Provider retry");
|
|
assert!(serialized.contains(PROVIDER_RETRY_SCHEMA_VERSION));
|
|
assert!(serialized.contains("\"baseRequestSlot\""));
|
|
assert!(serialized.contains("\"allowIdleContextCompaction\""));
|
|
|
|
let mut unexpected = serde_json::to_value(&second).expect("serialize Provider retry");
|
|
unexpected
|
|
.as_object_mut()
|
|
.expect("Provider retry object")
|
|
.insert("unexpected".to_string(), serde_json::json!(true));
|
|
fs::write(
|
|
provider_retry_path(directory.path(), &identity.agent_id, &identity.run_id),
|
|
serde_json::to_vec_pretty(&unexpected).expect("serialize unexpected Provider retry"),
|
|
)
|
|
.expect("write unexpected Provider retry");
|
|
let error = read_for_run_at(directory.path(), &identity.agent_id, &identity.run_id)
|
|
.expect_err("unknown Provider retry fields must fail");
|
|
assert!(error.contains("unknown field"));
|
|
}
|
|
|
|
#[test]
|
|
fn provider_retry_reads_atomic_previous_when_primary_is_missing() {
|
|
let directory = tempdir().expect("create temp directory");
|
|
let identity = identity("quality-lead", "run-previous");
|
|
let record = write_first(directory.path(), &identity);
|
|
let path = provider_retry_path(directory.path(), &identity.agent_id, &identity.run_id);
|
|
let backup_path = agent_runtime_json_sidecar_backup_path(&path);
|
|
fs::rename(&path, &backup_path).expect("move Provider retry to previous");
|
|
|
|
assert!(!path.exists());
|
|
assert_eq!(
|
|
read_for_run_at(directory.path(), &identity.agent_id, &identity.run_id)
|
|
.expect("read previous Provider retry"),
|
|
Some(record)
|
|
);
|
|
}
|
|
|
|
#[test]
|
|
fn provider_retry_list_sorts_across_agents_and_runs() {
|
|
let directory = tempdir().expect("create temp directory");
|
|
let identities = [
|
|
identity("quality-lead", "run-2"),
|
|
identity("design-director", "run-2"),
|
|
identity("quality-lead", "run-1"),
|
|
identity("design-director", "run-1"),
|
|
];
|
|
for identity in &identities {
|
|
write_first(directory.path(), identity);
|
|
}
|
|
|
|
let listed = list_at(directory.path()).expect("list sorted Provider retries");
|
|
let listed_keys = listed
|
|
.iter()
|
|
.map(|record| {
|
|
(
|
|
record.identity.agent_id.as_str(),
|
|
record.identity.run_id.as_str(),
|
|
)
|
|
})
|
|
.collect::<Vec<_>>();
|
|
assert_eq!(
|
|
listed_keys,
|
|
vec![
|
|
("design-director", "run-1"),
|
|
("design-director", "run-2"),
|
|
("quality-lead", "run-1"),
|
|
("quality-lead", "run-2"),
|
|
]
|
|
);
|
|
}
|
|
|
|
#[test]
|
|
fn provider_retry_list_recovers_when_only_atomic_previous_exists() {
|
|
let directory = tempdir().expect("create temp directory");
|
|
let identity = identity("quality-lead", "run-list-previous");
|
|
let record = write_first(directory.path(), &identity);
|
|
let path = provider_retry_path(directory.path(), &identity.agent_id, &identity.run_id);
|
|
let backup_path = agent_runtime_json_sidecar_backup_path(&path);
|
|
fs::rename(&path, &backup_path).expect("move Provider retry to previous");
|
|
|
|
assert_eq!(
|
|
list_at(directory.path()).expect("list previous Provider retry"),
|
|
vec![record]
|
|
);
|
|
}
|
|
|
|
#[test]
|
|
fn provider_retry_list_deduplicates_primary_and_atomic_previous() {
|
|
let directory = tempdir().expect("create temp directory");
|
|
let identity = identity("code-prototype", "run-list-deduplicate");
|
|
let record = write_first(directory.path(), &identity);
|
|
let path = provider_retry_path(directory.path(), &identity.agent_id, &identity.run_id);
|
|
let backup_path = agent_runtime_json_sidecar_backup_path(&path);
|
|
fs::copy(&path, &backup_path).expect("copy Provider retry to previous");
|
|
|
|
assert_eq!(
|
|
list_at(directory.path()).expect("list deduplicated Provider retry"),
|
|
vec![record]
|
|
);
|
|
}
|
|
|
|
#[test]
|
|
fn provider_retry_list_rejects_unknown_directory_entries() {
|
|
let directory = tempdir().expect("create temp directory");
|
|
let identity = identity("code-prototype", "run-list-unknown");
|
|
write_first(directory.path(), &identity);
|
|
let agent_directory =
|
|
provider_retry_path(directory.path(), &identity.agent_id, &identity.run_id)
|
|
.parent()
|
|
.expect("Provider retry Agent directory")
|
|
.to_path_buf();
|
|
fs::write(agent_directory.join("unexpected.txt"), b"unexpected")
|
|
.expect("write unknown Provider retry entry");
|
|
|
|
let error = list_at(directory.path()).expect_err("unknown entries must fail closed");
|
|
assert!(error.contains("Provider 重试目录包含未知文件:unexpected.txt"));
|
|
}
|
|
|
|
#[test]
|
|
fn provider_retry_rejects_identity_and_same_attempt_conflicts() {
|
|
let directory = tempdir().expect("create temp directory");
|
|
let identity = identity("code-prototype", "run-conflict");
|
|
write_first(directory.path(), &identity);
|
|
|
|
let mut conflicting_identity = identity.clone();
|
|
conflicting_identity.task_id = "other-task".to_string();
|
|
let identity_error = read_matching_at(directory.path(), &conflicting_identity)
|
|
.expect_err("identity conflict must fail");
|
|
assert!(identity_error.contains("身份冲突"));
|
|
let write_error = write_next_at(
|
|
directory.path(),
|
|
&conflicting_identity,
|
|
"loop-2-repair-0-transient-1",
|
|
1,
|
|
3,
|
|
250,
|
|
"transport",
|
|
&"b".repeat(64),
|
|
)
|
|
.expect_err("conflicting identity write must fail");
|
|
assert!(write_error.contains("身份冲突"));
|
|
|
|
let payload_error = write_next_at(
|
|
directory.path(),
|
|
&identity,
|
|
"loop-2-repair-0-transient-1",
|
|
1,
|
|
3,
|
|
250,
|
|
"connectivity",
|
|
&"c".repeat(64),
|
|
)
|
|
.expect_err("same attempt payload conflict must fail");
|
|
assert!(payload_error.contains("同一 attempt"));
|
|
}
|
|
|
|
#[test]
|
|
fn provider_retry_dangerous_path_components_do_not_collide() {
|
|
let first_agent = "design/director";
|
|
let second_agent = "design?director";
|
|
let run_id = "../run";
|
|
let first = provider_retry_relative_path(first_agent, run_id);
|
|
let second = provider_retry_relative_path(second_agent, run_id);
|
|
assert_ne!(first, second);
|
|
assert!(first.contains("design-director--"));
|
|
assert!(first.contains(&format!("{:x}", Sha256::digest(first_agent.as_bytes()))));
|
|
assert!(first.contains(&format!("{:x}", Sha256::digest(run_id.as_bytes()))));
|
|
assert!(!first.contains("/../"));
|
|
assert_eq!(
|
|
provider_retry_relative_path("code-prototype", "run-safe"),
|
|
".agent/runtime/provider-retries/code-prototype/run-safe.json"
|
|
);
|
|
assert_ne!(
|
|
provider_retry_path_component(".", "agent"),
|
|
provider_retry_path_component("..", "agent")
|
|
);
|
|
}
|
|
|
|
#[test]
|
|
fn provider_retry_remaining_time_has_strict_boundaries_and_can_be_forced_due() {
|
|
assert_eq!(remaining_ms_at(1_500, 1_000), 500);
|
|
assert_eq!(remaining_ms_at(1_500, 1_500), 0);
|
|
assert_eq!(remaining_ms_at(1_500, 2_000), 0);
|
|
|
|
let directory = tempdir().expect("create temp directory");
|
|
let identity = identity("design-director", "run-due");
|
|
write_first(directory.path(), &identity);
|
|
let due = force_provider_retry_due_for_test_at(directory.path(), &identity)
|
|
.expect("force Provider retry due");
|
|
assert_eq!(remaining_ms(&due), 0);
|
|
}
|
|
|
|
#[test]
|
|
fn provider_retry_remove_deletes_primary_and_previous_idempotently() {
|
|
let directory = tempdir().expect("create temp directory");
|
|
let identity = identity("quality-lead", "run-remove");
|
|
write_first(directory.path(), &identity);
|
|
let path = provider_retry_path(directory.path(), &identity.agent_id, &identity.run_id);
|
|
let backup_path = agent_runtime_json_sidecar_backup_path(&path);
|
|
fs::copy(&path, &backup_path).expect("copy Provider retry previous");
|
|
|
|
remove_at(directory.path(), &identity.agent_id, &identity.run_id)
|
|
.expect("remove Provider retry");
|
|
assert!(!path.exists());
|
|
assert!(!backup_path.exists());
|
|
remove_at(directory.path(), &identity.agent_id, &identity.run_id)
|
|
.expect("repeat Provider retry removal");
|
|
}
|
|
}
|