2477 lines
88 KiB
Rust
2477 lines
88 KiB
Rust
use super::*;
|
||
|
||
fn runtime(status: &str, phase: &str, pending: u32) -> AgentRuntimeResult {
|
||
let state = serde_json::from_value::<AgentRuntimeState>(serde_json::json!({
|
||
"agentId": "code-prototype",
|
||
"runId": "run-test",
|
||
"status": status,
|
||
"phase": phase,
|
||
}))
|
||
.expect("deserialize runtime fixture");
|
||
let mut task_queue = AgentRuntimeTaskQueueSummary::default();
|
||
task_queue.pending = pending;
|
||
AgentRuntimeResult {
|
||
state,
|
||
accepted_run_id: None,
|
||
session_path: String::new(),
|
||
event_path: String::new(),
|
||
task_path: String::new(),
|
||
task_queue,
|
||
recent_events: Vec::new(),
|
||
recent_tasks: Vec::new(),
|
||
response_stream: None,
|
||
user_input_request: None,
|
||
}
|
||
}
|
||
|
||
fn response_stream(
|
||
request_slot: &str,
|
||
response_revision: u64,
|
||
sequence: u64,
|
||
status: &str,
|
||
accumulated_text: &str,
|
||
) -> AgentRuntimeResponseStream {
|
||
AgentRuntimeResponseStream {
|
||
schema_version: "game-creator-runtime-response-stream.v1".to_string(),
|
||
agent_id: "code-prototype".to_string(),
|
||
task_id: "code-prototype".to_string(),
|
||
session_id: "session-test".to_string(),
|
||
run_id: "run-test".to_string(),
|
||
request_kind: "final-reply".to_string(),
|
||
request_slot: request_slot.to_string(),
|
||
applied_steer_cursor: 0,
|
||
response_revision,
|
||
sequence,
|
||
status: status.to_string(),
|
||
accumulated_text: accumulated_text.to_string(),
|
||
finish_reason: None,
|
||
started_at: 100,
|
||
updated_at: 100 + sequence,
|
||
}
|
||
}
|
||
|
||
fn runtime_with_response_stream(stream: AgentRuntimeResponseStream) -> AgentRuntimeResult {
|
||
let mut snapshot = runtime("running", "response", 0);
|
||
snapshot.state.session_id = stream.session_id.clone();
|
||
snapshot.response_stream = Some(stream);
|
||
snapshot
|
||
}
|
||
|
||
#[test]
|
||
fn new_supervisor_runs_fix_cli_source_and_select_requested_profile() {
|
||
for run_profile in [
|
||
AGENT_RUNTIME_RUN_PROFILE_STANDARD,
|
||
AGENT_RUNTIME_RUN_PROFILE_AUTONOMOUS_GAME_BUILD,
|
||
] {
|
||
assert_eq!(
|
||
resolve_swarm_new_run_launch(
|
||
GAME_CREATOR_PROJECT_SUPERVISOR_AGENT_ID,
|
||
run_profile,
|
||
AGENT_RUNTIME_SUPERVISOR_CLI_SOURCE,
|
||
)
|
||
.expect("resolve supervisor launch"),
|
||
SwarmNewRunLaunch::ProjectSupervisor {
|
||
source: AGENT_RUNTIME_SUPERVISOR_CLI_SOURCE,
|
||
run_profile,
|
||
}
|
||
);
|
||
}
|
||
}
|
||
|
||
#[test]
|
||
fn explicit_parent_debug_keeps_standard_profile_only() {
|
||
assert_eq!(
|
||
resolve_swarm_new_run_launch(
|
||
"code-prototype",
|
||
AGENT_RUNTIME_RUN_PROFILE_STANDARD,
|
||
AGENT_RUNTIME_SUPERVISOR_CLI_SOURCE,
|
||
)
|
||
.expect("resolve explicit parent debug launch"),
|
||
SwarmNewRunLaunch::ExplicitParentDebug,
|
||
);
|
||
assert!(resolve_swarm_new_run_launch(
|
||
"code-prototype",
|
||
AGENT_RUNTIME_RUN_PROFILE_AUTONOMOUS_GAME_BUILD,
|
||
AGENT_RUNTIME_SUPERVISOR_CLI_SOURCE,
|
||
)
|
||
.expect_err("autonomous profile must stay on the supervisor root run")
|
||
.contains(GAME_CREATOR_PROJECT_SUPERVISOR_AGENT_ID));
|
||
assert!(resolve_swarm_new_run_launch(
|
||
GAME_CREATOR_PROJECT_SUPERVISOR_AGENT_ID,
|
||
"unsupported",
|
||
AGENT_RUNTIME_SUPERVISOR_CLI_SOURCE,
|
||
)
|
||
.is_err());
|
||
}
|
||
|
||
#[test]
|
||
fn plan_entry_never_defers_to_the_interaction_kernel() {
|
||
// GUI 的「做方案」是一个直接起 plan 根 Run 的按钮;无头入口若把 reply/execute
|
||
// 的判定交给模型,同一句需求就会时而起 Run、时而只回一段口头建议。
|
||
assert!(!swarm_turn_uses_interaction_kernel(
|
||
GAME_CREATOR_PROJECT_SUPERVISOR_AGENT_ID,
|
||
Some(AGENT_RUNTIME_SUPERVISOR_PLAN_SOURCE),
|
||
));
|
||
// 做游戏 / 做素材两条链路继续走交互内核,行为不变。
|
||
assert_eq!(
|
||
swarm_turn_uses_interaction_kernel(GAME_CREATOR_PROJECT_SUPERVISOR_AGENT_ID, None),
|
||
game_creator_agent_uses_interaction_kernel(GAME_CREATOR_PROJECT_SUPERVISOR_AGENT_ID)
|
||
);
|
||
assert_eq!(
|
||
swarm_turn_uses_interaction_kernel(
|
||
GAME_CREATOR_PROJECT_SUPERVISOR_AGENT_ID,
|
||
Some(AGENT_RUNTIME_SUPERVISOR_CLI_SOURCE),
|
||
),
|
||
game_creator_agent_uses_interaction_kernel(GAME_CREATOR_PROJECT_SUPERVISOR_AGENT_ID)
|
||
);
|
||
}
|
||
|
||
#[test]
|
||
fn plan_launch_claims_its_own_source_and_rejects_autonomous_profile() {
|
||
assert_eq!(
|
||
resolve_swarm_new_run_launch(
|
||
GAME_CREATOR_PROJECT_SUPERVISOR_AGENT_ID,
|
||
AGENT_RUNTIME_RUN_PROFILE_STANDARD,
|
||
AGENT_RUNTIME_SUPERVISOR_PLAN_SOURCE,
|
||
)
|
||
.expect("resolve plan launch"),
|
||
SwarmNewRunLaunch::ProjectSupervisor {
|
||
source: AGENT_RUNTIME_SUPERVISOR_PLAN_SOURCE,
|
||
run_profile: AGENT_RUNTIME_RUN_PROFILE_STANDARD,
|
||
}
|
||
);
|
||
// 立项策划 Run 和做游戏 Run 都是 standard 档,只有按 source 认领才不会串链。
|
||
assert_eq!(
|
||
SwarmNewRunLaunch::ProjectSupervisor {
|
||
source: AGENT_RUNTIME_SUPERVISOR_PLAN_SOURCE,
|
||
run_profile: AGENT_RUNTIME_RUN_PROFILE_STANDARD,
|
||
}
|
||
.expected_parent_source(),
|
||
Some(AGENT_RUNTIME_SUPERVISOR_PLAN_SOURCE)
|
||
);
|
||
assert!(resolve_swarm_new_run_launch(
|
||
GAME_CREATOR_PROJECT_SUPERVISOR_AGENT_ID,
|
||
AGENT_RUNTIME_RUN_PROFILE_AUTONOMOUS_GAME_BUILD,
|
||
AGENT_RUNTIME_SUPERVISOR_PLAN_SOURCE,
|
||
)
|
||
.is_err());
|
||
assert!(resolve_swarm_new_run_launch(
|
||
"code-prototype",
|
||
AGENT_RUNTIME_RUN_PROFILE_STANDARD,
|
||
AGENT_RUNTIME_SUPERVISOR_PLAN_SOURCE,
|
||
)
|
||
.is_err());
|
||
}
|
||
|
||
#[test]
|
||
fn same_run_steer_preserves_bound_autonomous_profile() {
|
||
let root = std::env::temp_dir().join(format!(
|
||
"swarm-cli-profile-steer-{}-{}",
|
||
std::process::id(),
|
||
unix_millis()
|
||
));
|
||
init_local_game_project_at(&root, "project-swarm-profile", "Swarm Profile Steer")
|
||
.expect("initialize profile steer project");
|
||
let run_id = "swarm-profile-steer-run";
|
||
let binding = bind_game_creator_agent_runtime_run_profile_at(
|
||
&root,
|
||
GAME_CREATOR_PROJECT_SUPERVISOR_AGENT_ID,
|
||
run_id,
|
||
AGENT_RUNTIME_SUPERVISOR_CLI_SOURCE,
|
||
Some(AGENT_RUNTIME_RUN_PROFILE_AUTONOMOUS_GAME_BUILD),
|
||
None,
|
||
)
|
||
.expect("bind autonomous supervisor profile");
|
||
let state = start_game_creator_agent_runtime_task_at(
|
||
&root,
|
||
GAME_CREATOR_PROJECT_SUPERVISOR_AGENT_ID,
|
||
"生成一版可试玩项目",
|
||
run_id,
|
||
AGENT_RUNTIME_SUPERVISOR_CLI_SOURCE,
|
||
"准备自主构建",
|
||
vec!["实现并验证最小可玩闭环".to_string()],
|
||
)
|
||
.expect("start autonomous supervisor runtime");
|
||
|
||
let steered = steer_game_creator_agent_runtime_task_at(
|
||
&root,
|
||
GAME_CREATOR_PROJECT_SUPERVISOR_AGENT_ID,
|
||
&state.session_id,
|
||
run_id,
|
||
"swarm-profile-steer-1",
|
||
"保持当前目标并补充触屏操作",
|
||
"swarm-cli",
|
||
)
|
||
.expect("steer autonomous supervisor runtime");
|
||
|
||
assert_eq!(steered.runtime.state.run_id, run_id);
|
||
assert_eq!(
|
||
steered.runtime.state.run_profile,
|
||
AGENT_RUNTIME_RUN_PROFILE_AUTONOMOUS_GAME_BUILD
|
||
);
|
||
assert_eq!(
|
||
steered.runtime.state.run_profile_binding_fingerprint,
|
||
binding.binding_fingerprint
|
||
);
|
||
fs::remove_dir_all(root).ok();
|
||
}
|
||
|
||
#[test]
|
||
fn cross_profile_steer_is_rejected_before_persistent_side_effects() {
|
||
let root = std::env::temp_dir().join(format!(
|
||
"swarm-cli-profile-mismatch-{}-{}",
|
||
std::process::id(),
|
||
unix_millis()
|
||
));
|
||
init_local_game_project_at(&root, "project-swarm-mismatch", "Swarm Profile Mismatch")
|
||
.expect("initialize profile mismatch project");
|
||
let run_id = "swarm-profile-standard-run";
|
||
let state = start_game_creator_agent_runtime_task_at(
|
||
&root,
|
||
GAME_CREATOR_PROJECT_SUPERVISOR_AGENT_ID,
|
||
"等待开发者确认的标准任务",
|
||
run_id,
|
||
AGENT_RUNTIME_SUPERVISOR_CLI_SOURCE,
|
||
"准备标准模式运行",
|
||
vec!["等待确认".to_string()],
|
||
)
|
||
.expect("start standard supervisor runtime");
|
||
assert_eq!(
|
||
state.run_profile, AGENT_RUNTIME_RUN_PROFILE_STANDARD,
|
||
"fixture must stay standard"
|
||
);
|
||
let conversation_before = read_local_conversation_for_session_at(
|
||
&root,
|
||
Some(GAME_CREATOR_PROJECT_SUPERVISOR_AGENT_ID),
|
||
Some(&state.session_id),
|
||
)
|
||
.expect("read conversation before rejected steer");
|
||
|
||
let error = steer_game_creator_agent_runtime_task_for_profile_at(
|
||
&root,
|
||
GAME_CREATOR_PROJECT_SUPERVISOR_AGENT_ID,
|
||
&state.session_id,
|
||
run_id,
|
||
"swarm-profile-mismatch-steer",
|
||
"切换为自主构建并继续",
|
||
Some(AGENT_RUNTIME_RUN_PROFILE_AUTONOMOUS_GAME_BUILD),
|
||
"swarm-cli",
|
||
)
|
||
.expect_err("cross-profile steer must fail closed");
|
||
assert!(error.contains("Run Profile"), "unexpected error: {error}");
|
||
assert!(!game_creator_agent_runtime_steer_ledger_path(
|
||
&root,
|
||
GAME_CREATOR_PROJECT_SUPERVISOR_AGENT_ID,
|
||
run_id,
|
||
)
|
||
.exists());
|
||
let conversation_after = read_local_conversation_for_session_at(
|
||
&root,
|
||
Some(GAME_CREATOR_PROJECT_SUPERVISOR_AGENT_ID),
|
||
Some(&state.session_id),
|
||
)
|
||
.expect("read conversation after rejected steer");
|
||
assert_eq!(conversation_after.messages, conversation_before.messages);
|
||
let persisted =
|
||
read_game_creator_agent_runtime_at(&root, GAME_CREATOR_PROJECT_SUPERVISOR_AGENT_ID)
|
||
.expect("read runtime after rejected steer");
|
||
assert_eq!(persisted.state.run_id, run_id);
|
||
assert_eq!(
|
||
persisted.state.run_profile,
|
||
AGENT_RUNTIME_RUN_PROFILE_STANDARD
|
||
);
|
||
fs::remove_dir_all(root).ok();
|
||
}
|
||
|
||
#[test]
|
||
fn natural_language_is_not_preclassified_by_cli_keywords() {
|
||
for message in [
|
||
"你是谁",
|
||
"项目在哪里",
|
||
"实现登录页并运行测试",
|
||
"不要运行任何命令,只解释架构",
|
||
"你能做什么,并顺便修复这个问题",
|
||
"继续刚才的任务",
|
||
] {
|
||
assert_eq!(
|
||
parse_swarm_chat_input(message),
|
||
Some(SwarmChatInput::Message(message.to_string())),
|
||
"{message} must enter the unified Agent interaction loop"
|
||
);
|
||
}
|
||
}
|
||
|
||
#[test]
|
||
fn resume_is_an_explicit_control_command() {
|
||
assert_eq!(
|
||
parse_swarm_chat_input("/resume"),
|
||
Some(SwarmChatInput::Resume)
|
||
);
|
||
}
|
||
|
||
#[test]
|
||
fn natural_language_resume_without_active_runtime_starts_a_new_run() {
|
||
assert_eq!(
|
||
normalize_interaction_action_without_active_runtime(AgentInteractionAction::Resume),
|
||
AgentInteractionAction::Execute
|
||
);
|
||
assert_eq!(
|
||
normalize_interaction_action_without_active_runtime(AgentInteractionAction::Reply(
|
||
"记得之前的工作".to_string()
|
||
)),
|
||
AgentInteractionAction::Reply("记得之前的工作".to_string())
|
||
);
|
||
}
|
||
|
||
#[test]
|
||
fn parent_runtime_matching_is_scoped_to_requested_profile() {
|
||
let mut standard = runtime("running", "planning", 0);
|
||
standard.state.agent_id = GAME_CREATOR_PROJECT_SUPERVISOR_AGENT_ID.to_string();
|
||
standard.state.session_id = "session-profile-match".to_string();
|
||
standard.state.run_id = "run-standard".to_string();
|
||
standard.state.run_profile = AGENT_RUNTIME_RUN_PROFILE_STANDARD.to_string();
|
||
let mut autonomous = standard.clone();
|
||
autonomous.state.run_id = "run-autonomous".to_string();
|
||
autonomous.state.run_profile = AGENT_RUNTIME_RUN_PROFILE_AUTONOMOUS_GAME_BUILD.to_string();
|
||
let runtimes = vec![standard, autonomous];
|
||
|
||
assert_eq!(
|
||
swarm_parent_runtime(
|
||
GAME_CREATOR_PROJECT_SUPERVISOR_AGENT_ID,
|
||
"session-profile-match",
|
||
AGENT_RUNTIME_RUN_PROFILE_STANDARD,
|
||
None,
|
||
&runtimes,
|
||
)
|
||
.map(|runtime| runtime.state.run_id.as_str()),
|
||
Some("run-standard")
|
||
);
|
||
assert_eq!(
|
||
swarm_parent_runtime(
|
||
GAME_CREATOR_PROJECT_SUPERVISOR_AGENT_ID,
|
||
"session-profile-match",
|
||
AGENT_RUNTIME_RUN_PROFILE_AUTONOMOUS_GAME_BUILD,
|
||
None,
|
||
&runtimes,
|
||
)
|
||
.map(|runtime| runtime.state.run_id.as_str()),
|
||
Some("run-autonomous")
|
||
);
|
||
}
|
||
|
||
#[test]
|
||
fn completed_parent_with_pending_task_is_busy_but_not_a_steer_target() {
|
||
let session_id = "session-pending-after-completed";
|
||
let run_profile = AGENT_RUNTIME_RUN_PROFILE_AUTONOMOUS_GAME_BUILD;
|
||
let mut completed = runtime("idle", "completed", 1);
|
||
completed.state.agent_id = GAME_CREATOR_PROJECT_SUPERVISOR_AGENT_ID.to_string();
|
||
completed.state.session_id = session_id.to_string();
|
||
completed.state.run_id = "run-completed-before-pending".to_string();
|
||
completed.state.run_profile = run_profile.to_string();
|
||
completed.recent_tasks.push(
|
||
serde_json::from_value(serde_json::json!({
|
||
"agentId": GAME_CREATOR_PROJECT_SUPERVISOR_AGENT_ID,
|
||
"taskId": GAME_CREATOR_PROJECT_SUPERVISOR_AGENT_ID,
|
||
"sessionId": session_id,
|
||
"runId": "run-pending-after-completed",
|
||
"source": AGENT_RUNTIME_SUPERVISOR_CLI_SOURCE,
|
||
"runProfile": run_profile,
|
||
"task": "继续补齐游戏功能",
|
||
"status": "pending",
|
||
"phase": "queued"
|
||
}))
|
||
.expect("deserialize pending task fixture"),
|
||
);
|
||
let runtimes = vec![completed];
|
||
|
||
assert!(runtime_is_busy(&runtimes[0]));
|
||
assert!(
|
||
swarm_parent_steer_target(
|
||
GAME_CREATOR_PROJECT_SUPERVISOR_AGENT_ID,
|
||
session_id,
|
||
run_profile,
|
||
None,
|
||
None,
|
||
&runtimes,
|
||
)
|
||
.is_none(),
|
||
"queued work must not make the completed canonical run steerable"
|
||
);
|
||
assert_eq!(
|
||
matching_pending_swarm_run_id(
|
||
&runtimes[0],
|
||
session_id,
|
||
run_profile,
|
||
None,
|
||
"继续补齐游戏功能",
|
||
),
|
||
Some("run-pending-after-completed")
|
||
);
|
||
}
|
||
|
||
#[test]
|
||
fn new_turn_uses_accepted_run_id_instead_of_stale_canonical_state() {
|
||
let mut started = runtime("cancelled", "cancelled", 1);
|
||
started.state.agent_id = GAME_CREATOR_PROJECT_SUPERVISOR_AGENT_ID.to_string();
|
||
started.state.run_id = "run-old-cancelled".to_string();
|
||
started.accepted_run_id = Some("run-new-pending".to_string());
|
||
|
||
assert_eq!(
|
||
accepted_swarm_run_id(&started, "run-requested"),
|
||
"run-new-pending"
|
||
);
|
||
started.accepted_run_id = None;
|
||
assert_eq!(
|
||
accepted_swarm_run_id(&started, "run-requested"),
|
||
"run-requested"
|
||
);
|
||
}
|
||
|
||
#[test]
|
||
fn queued_start_returns_actual_accepted_run_id_after_collision() {
|
||
let root = std::env::temp_dir().join(format!(
|
||
"swarm-cli-accepted-run-id-{}-{}",
|
||
std::process::id(),
|
||
unix_millis()
|
||
));
|
||
init_local_game_project_at(&root, "project-accepted-run-id", "Accepted run ID")
|
||
.expect("initialize accepted run ID project");
|
||
let runtime_lock = try_acquire_game_creator_agent_runtime_task_lock(
|
||
&root,
|
||
GAME_CREATOR_PROJECT_SUPERVISOR_AGENT_ID,
|
||
)
|
||
.expect("acquire supervisor runtime lock")
|
||
.expect("supervisor runtime lock available");
|
||
|
||
let first = start_game_creator_supervisor_background_task_for_session_at(
|
||
&root,
|
||
None,
|
||
"第一条排队任务",
|
||
"run-collision",
|
||
AGENT_RUNTIME_SUPERVISOR_CLI_SOURCE,
|
||
AGENT_RUNTIME_RUN_PROFILE_AUTONOMOUS_GAME_BUILD,
|
||
)
|
||
.expect("queue first colliding run");
|
||
let second = start_game_creator_supervisor_background_task_for_session_at(
|
||
&root,
|
||
None,
|
||
"第二条排队任务",
|
||
"run-collision",
|
||
AGENT_RUNTIME_SUPERVISOR_CLI_SOURCE,
|
||
AGENT_RUNTIME_RUN_PROFILE_AUTONOMOUS_GAME_BUILD,
|
||
)
|
||
.expect("queue second colliding run");
|
||
|
||
assert_eq!(first.accepted_run_id.as_deref(), Some("run-collision"));
|
||
let second_run_id = second
|
||
.accepted_run_id
|
||
.as_deref()
|
||
.expect("second accepted run ID");
|
||
assert_ne!(second_run_id, "run-collision");
|
||
assert!(second_run_id.starts_with("run-collision-dup-"));
|
||
let serialized = serde_json::to_value(&second).expect("serialize queued start result");
|
||
assert_eq!(serialized["acceptedRunId"], second_run_id);
|
||
assert!(serialized.get("accepted_run_id").is_none());
|
||
|
||
drop(runtime_lock);
|
||
fs::remove_dir_all(root).ok();
|
||
}
|
||
|
||
#[test]
|
||
fn completed_target_turn_settles_after_canonical_advances_to_next_run() {
|
||
let root = std::env::temp_dir().join(format!(
|
||
"swarm-cli-consecutive-snapshot-{}-{}",
|
||
std::process::id(),
|
||
unix_millis()
|
||
));
|
||
init_local_game_project_at(
|
||
&root,
|
||
"project-consecutive-snapshot",
|
||
"Consecutive run snapshot",
|
||
)
|
||
.expect("initialize consecutive snapshot project");
|
||
let session_id = "session-consecutive-runs";
|
||
let run_profile = AGENT_RUNTIME_RUN_PROFILE_AUTONOMOUS_GAME_BUILD;
|
||
let completed_task = serde_json::from_value::<AgentRuntimeTaskRecord>(serde_json::json!({
|
||
"agentId": GAME_CREATOR_PROJECT_SUPERVISOR_AGENT_ID,
|
||
"taskId": GAME_CREATOR_PROJECT_SUPERVISOR_AGENT_ID,
|
||
"sessionId": session_id,
|
||
"runId": "run-target-completed",
|
||
"source": AGENT_RUNTIME_SUPERVISOR_CLI_SOURCE,
|
||
"runProfile": run_profile,
|
||
"task": "先完成这一轮",
|
||
"status": "completed",
|
||
"phase": "completed",
|
||
"currentAction": "本轮已经完成"
|
||
}))
|
||
.expect("deserialize completed target task");
|
||
let next_task = serde_json::from_value::<AgentRuntimeTaskRecord>(serde_json::json!({
|
||
"agentId": GAME_CREATOR_PROJECT_SUPERVISOR_AGENT_ID,
|
||
"taskId": GAME_CREATOR_PROJECT_SUPERVISOR_AGENT_ID,
|
||
"sessionId": session_id,
|
||
"runId": "run-next-running",
|
||
"source": AGENT_RUNTIME_SUPERVISOR_CLI_SOURCE,
|
||
"runProfile": run_profile,
|
||
"task": "随后执行下一轮",
|
||
"status": "running",
|
||
"phase": "planning",
|
||
"currentAction": "下一轮正在执行"
|
||
}))
|
||
.expect("deserialize next running task");
|
||
let mut current = runtime("running", "planning", 0);
|
||
current.state.agent_id = GAME_CREATOR_PROJECT_SUPERVISOR_AGENT_ID.to_string();
|
||
current.state.session_id = session_id.to_string();
|
||
current.state.run_id = "run-next-running".to_string();
|
||
current.state.run_profile = run_profile.to_string();
|
||
let task_path =
|
||
game_creator_agent_runtime_task_path(&root, GAME_CREATOR_PROJECT_SUPERVISOR_AGENT_ID);
|
||
fs::create_dir_all(task_path.parent().expect("task journal parent"))
|
||
.expect("create task journal parent");
|
||
fs::write(
|
||
&task_path,
|
||
format!(
|
||
"{}\n{}\n",
|
||
serde_json::to_string(&completed_task).expect("serialize completed target task"),
|
||
serde_json::to_string(&next_task).expect("serialize next target task"),
|
||
),
|
||
)
|
||
.expect("persist full task journal");
|
||
current.recent_tasks = vec![next_task];
|
||
let runtimes = vec![current];
|
||
|
||
assert!(!swarm_turn_is_busy(
|
||
&root,
|
||
GAME_CREATOR_PROJECT_SUPERVISOR_AGENT_ID,
|
||
session_id,
|
||
run_profile,
|
||
None,
|
||
"run-target-completed",
|
||
&runtimes,
|
||
)
|
||
.expect("read completed target busy state"));
|
||
assert!(swarm_turn_is_busy(
|
||
&root,
|
||
GAME_CREATOR_PROJECT_SUPERVISOR_AGENT_ID,
|
||
session_id,
|
||
run_profile,
|
||
None,
|
||
"run-next-running",
|
||
&runtimes,
|
||
)
|
||
.expect("read next target busy state"));
|
||
let snapshot = swarm_parent_runtime_snapshot_for_run(
|
||
&root,
|
||
GAME_CREATOR_PROJECT_SUPERVISOR_AGENT_ID,
|
||
session_id,
|
||
run_profile,
|
||
None,
|
||
"run-target-completed",
|
||
&runtimes,
|
||
)
|
||
.expect("read target task journal")
|
||
.expect("recover completed target from task journal");
|
||
assert!(parent_runtime_completed(&snapshot));
|
||
|
||
fs::remove_dir_all(root).ok();
|
||
}
|
||
|
||
#[test]
|
||
fn running_parent_remains_the_only_valid_swarm_steer_target() {
|
||
let session_id = "session-running-steer";
|
||
let run_profile = AGENT_RUNTIME_RUN_PROFILE_STANDARD;
|
||
let mut running = runtime("running", "waiting-for-provider-retry", 0);
|
||
running.state.agent_id = GAME_CREATOR_PROJECT_SUPERVISOR_AGENT_ID.to_string();
|
||
running.state.session_id = session_id.to_string();
|
||
running.state.run_id = "run-running-steer".to_string();
|
||
running.state.run_profile = run_profile.to_string();
|
||
let runtimes = vec![running];
|
||
|
||
assert_eq!(
|
||
swarm_parent_steer_target(
|
||
GAME_CREATOR_PROJECT_SUPERVISOR_AGENT_ID,
|
||
session_id,
|
||
run_profile,
|
||
None,
|
||
Some("run-running-steer"),
|
||
&runtimes,
|
||
)
|
||
.map(|runtime| runtime.state.run_id.as_str()),
|
||
Some("run-running-steer")
|
||
);
|
||
assert!(swarm_parent_steer_target(
|
||
GAME_CREATOR_PROJECT_SUPERVISOR_AGENT_ID,
|
||
session_id,
|
||
run_profile,
|
||
None,
|
||
Some("another-run"),
|
||
&runtimes,
|
||
)
|
||
.is_none());
|
||
}
|
||
|
||
#[test]
|
||
fn parses_chat_commands_without_stealing_normal_messages() {
|
||
assert_eq!(parse_swarm_chat_input(" "), None);
|
||
assert_eq!(
|
||
parse_swarm_chat_input("/agents"),
|
||
Some(SwarmChatInput::Agents)
|
||
);
|
||
assert_eq!(parse_swarm_chat_input("/exit"), Some(SwarmChatInput::Quit));
|
||
assert_eq!(
|
||
parse_swarm_chat_input("/status"),
|
||
Some(SwarmChatInput::Status)
|
||
);
|
||
assert_eq!(
|
||
parse_swarm_chat_input("/compact"),
|
||
Some(SwarmChatInput::Compact)
|
||
);
|
||
assert_eq!(
|
||
parse_swarm_chat_input("让策划和程序并行检查玩法"),
|
||
Some(SwarmChatInput::Message(
|
||
"让策划和程序并行检查玩法".to_string()
|
||
))
|
||
);
|
||
}
|
||
|
||
#[test]
|
||
fn parses_goal_commands_and_keeps_goal_namespace_out_of_messages() {
|
||
assert_eq!(
|
||
parse_swarm_chat_input("/goal"),
|
||
Some(SwarmChatInput::Goal(SwarmGoalCommand::Status))
|
||
);
|
||
assert_eq!(
|
||
parse_swarm_chat_input("/goal status"),
|
||
Some(SwarmChatInput::Goal(SwarmGoalCommand::Status))
|
||
);
|
||
assert_eq!(
|
||
parse_swarm_chat_input("/goal 完成可玩的战斗循环"),
|
||
Some(SwarmChatInput::Goal(SwarmGoalCommand::Start(
|
||
"完成可玩的战斗循环".to_string()
|
||
)))
|
||
);
|
||
assert_eq!(
|
||
parse_swarm_chat_input("/goal edit 增加键盘与触屏验收"),
|
||
Some(SwarmChatInput::Goal(SwarmGoalCommand::Edit(
|
||
"增加键盘与触屏验收".to_string()
|
||
)))
|
||
);
|
||
assert_eq!(
|
||
parse_swarm_chat_input("/goal pause"),
|
||
Some(SwarmChatInput::Goal(SwarmGoalCommand::Pause))
|
||
);
|
||
assert_eq!(
|
||
parse_swarm_chat_input("/goal resume"),
|
||
Some(SwarmChatInput::Goal(SwarmGoalCommand::Resume))
|
||
);
|
||
assert_eq!(
|
||
parse_swarm_chat_input("/goal clear"),
|
||
Some(SwarmChatInput::Goal(SwarmGoalCommand::Clear))
|
||
);
|
||
|
||
for invalid in [
|
||
"/goal-status",
|
||
"/goal/status",
|
||
"/goal edit",
|
||
"/goal pause now",
|
||
] {
|
||
assert!(matches!(
|
||
parse_swarm_chat_input(invalid),
|
||
Some(SwarmChatInput::InvalidGoal(_))
|
||
));
|
||
assert!(!matches!(
|
||
parse_swarm_chat_input(invalid),
|
||
Some(SwarmChatInput::Message(_))
|
||
));
|
||
}
|
||
}
|
||
|
||
#[test]
|
||
fn swarm_help_lists_the_complete_goal_control_surface() {
|
||
let mut output = Vec::new();
|
||
print_swarm_chat_help(&mut output).expect("print swarm help");
|
||
let output = String::from_utf8(output).expect("help output is utf-8");
|
||
|
||
for command in [
|
||
"/status",
|
||
"/compact",
|
||
"/goal <目标>",
|
||
"/goal status",
|
||
"/goal edit <目标>",
|
||
"/goal pause",
|
||
"/goal resume",
|
||
"/goal clear",
|
||
] {
|
||
assert!(output.contains(command), "missing help command: {command}");
|
||
}
|
||
}
|
||
|
||
#[test]
|
||
fn goal_status_prints_identity_outcome_and_completion_standard() {
|
||
let goal = AgentGoalRecord {
|
||
schema_version: AGENT_GOAL_SCHEMA_VERSION.to_string(),
|
||
project_id: "project-1".to_string(),
|
||
goal_id: "goal-1".to_string(),
|
||
agent_id: "project-supervisor".to_string(),
|
||
session_id: "session-1".to_string(),
|
||
run_id: "run-1".to_string(),
|
||
revision: 3,
|
||
status: AGENT_GOAL_STATUS_ACTIVE.to_string(),
|
||
outcome: "完成首个可玩版本".to_string(),
|
||
constraints: vec!["不新增平行 Runtime".to_string()],
|
||
verification: vec!["键盘与触屏均可完成一局".to_string()],
|
||
completion_evidence: Vec::new(),
|
||
response_fingerprint: None,
|
||
created_at: 1,
|
||
pause_requested_at: None,
|
||
paused_at: None,
|
||
completed_at: None,
|
||
cleared_at: None,
|
||
error: None,
|
||
updated_at: 2,
|
||
};
|
||
let mut output = Vec::new();
|
||
print_swarm_goal_status("session-1", Some(&goal), &mut output).expect("print goal status");
|
||
let output = String::from_utf8(output).expect("goal output is utf-8");
|
||
|
||
assert!(output.contains("goal=goal-1 run=run-1 revision=3 status=active"));
|
||
assert!(output.contains("[Goal 目标] 完成首个可玩版本"));
|
||
assert!(output.contains("[Goal 约束] 不新增平行 Runtime"));
|
||
assert!(output.contains("[Goal 完成标准] 键盘与触屏均可完成一局"));
|
||
}
|
||
|
||
#[test]
|
||
fn swarm_stays_busy_for_active_queue_and_reconciliation() {
|
||
assert!(runtimes_are_busy(&[runtime("running", "planning", 0)]));
|
||
assert!(runtimes_are_busy(&[runtime("idle", "completed", 1)]));
|
||
assert!(runtimes_are_busy(&[runtime(
|
||
"failed",
|
||
"needs-reconciliation",
|
||
0
|
||
)]));
|
||
assert!(!runtimes_are_busy(&[runtime("idle", "completed", 0)]));
|
||
assert!(!runtimes_are_busy(&[runtime("failed", "failed", 0)]));
|
||
}
|
||
|
||
#[test]
|
||
fn active_turn_eof_closes_input_once_without_requesting_quit() {
|
||
let mut input_closed = false;
|
||
let mut observer = SwarmRuntimeObserver::default();
|
||
let mut output = Vec::new();
|
||
|
||
mark_swarm_turn_input_closed(&mut input_closed, &mut observer, &mut output)
|
||
.expect("close active turn input");
|
||
mark_swarm_turn_input_closed(&mut input_closed, &mut observer, &mut output)
|
||
.expect("repeat closed input is idempotent");
|
||
|
||
assert!(input_closed);
|
||
let output = String::from_utf8(output).expect("input close output is utf-8");
|
||
assert_eq!(output.matches("[输入已关闭]").count(), 1);
|
||
assert!(output.contains("继续运行,等待可信终态"));
|
||
assert!(!output.contains("已退出 Agent Swarm Chat"));
|
||
}
|
||
|
||
#[test]
|
||
fn active_turn_eof_keeps_observing_until_parent_completes() {
|
||
let root = std::env::temp_dir().join(format!(
|
||
"swarm-cli-eof-active-{}-{}",
|
||
std::process::id(),
|
||
unix_millis()
|
||
));
|
||
init_local_game_project_at(&root, "project-swarm-eof", "Swarm EOF active turn")
|
||
.expect("initialize EOF project");
|
||
for group in GAME_CREATOR_AGENT_GROUP_DEFINITIONS {
|
||
for role in group.roles {
|
||
let idle = default_game_creator_agent_runtime_state(role.task_id, "run-eof-idle");
|
||
write_game_creator_agent_runtime_state(&root, &idle)
|
||
.expect("persist valid idle specialist state");
|
||
}
|
||
}
|
||
let parent_agent_id = GAME_CREATOR_PROJECT_SUPERVISOR_AGENT_ID;
|
||
let before = append_local_conversation_message_at(
|
||
&root,
|
||
Some(parent_agent_id),
|
||
LocalConversationMessage {
|
||
role: "user".to_string(),
|
||
content: "继续完成当前项目".to_string(),
|
||
agent_id: Some(parent_agent_id.to_string()),
|
||
},
|
||
)
|
||
.expect("append turn user message");
|
||
let session_id = before.session_id.clone().expect("active parent session");
|
||
let mut parent = runtime("running", "planning", 0).state;
|
||
parent.agent_id = parent_agent_id.to_string();
|
||
parent.task_id = parent_agent_id.to_string();
|
||
parent.session_id = session_id.clone();
|
||
parent.run_id = "run-eof-active".to_string();
|
||
parent.source = "agent-background-task".to_string();
|
||
parent.current_task = "继续完成当前项目".to_string();
|
||
write_game_creator_agent_runtime_state(&root, &parent).expect("persist active parent");
|
||
let conversation_baseline =
|
||
new_swarm_turn_conversation_baseline(before.messages.len(), &parent.run_id);
|
||
|
||
let completion_root = root.clone();
|
||
let completion_session_id = session_id.clone();
|
||
let completion = std::thread::spawn(move || {
|
||
std::thread::sleep(Duration::from_millis(15));
|
||
append_local_conversation_message_for_session_at(
|
||
&completion_root,
|
||
Some(parent_agent_id),
|
||
Some(&completion_session_id),
|
||
LocalConversationMessage {
|
||
role: "assistant".to_string(),
|
||
content: "已完成可信终态".to_string(),
|
||
agent_id: Some(parent_agent_id.to_string()),
|
||
},
|
||
)
|
||
.expect("append terminal assistant message");
|
||
parent.status = "idle".to_string();
|
||
parent.phase = "completed".to_string();
|
||
parent.updated_at = unix_timestamp();
|
||
write_game_creator_agent_runtime_state(&completion_root, &parent)
|
||
.expect("persist completed parent");
|
||
});
|
||
let (tx, rx) = mpsc::channel();
|
||
tx.send(SwarmInputEvent::Eof).expect("send active turn EOF");
|
||
let mut observer = SwarmRuntimeObserver::default();
|
||
let mut output = Vec::new();
|
||
|
||
let outcome = wait_for_swarm_turn(
|
||
&root,
|
||
parent_agent_id,
|
||
&session_id,
|
||
AGENT_RUNTIME_RUN_PROFILE_STANDARD,
|
||
None,
|
||
conversation_baseline,
|
||
&rx,
|
||
&mut output,
|
||
&mut observer,
|
||
Duration::from_millis(2),
|
||
Duration::from_millis(8),
|
||
)
|
||
.expect("observe active turn after EOF");
|
||
completion.join().expect("join completion writer");
|
||
|
||
let output = String::from_utf8(output).expect("EOF turn output is utf-8");
|
||
let runtime_diagnostics = read_game_creator_agent_runtimes_at(&root)
|
||
.expect("read terminal runtime diagnostics")
|
||
.into_iter()
|
||
.filter(|runtime| runtime.state.phase == "needs-reconciliation")
|
||
.map(|runtime| {
|
||
format!(
|
||
"{}:{}",
|
||
runtime.state.agent_id,
|
||
runtime.state.error.unwrap_or_default()
|
||
)
|
||
})
|
||
.collect::<Vec<_>>();
|
||
assert!(
|
||
matches!(outcome, SwarmTurnOutcome::Settled(_)),
|
||
"unexpected outcome: {outcome:?}; diagnostics={runtime_diagnostics:?}; output={output}"
|
||
);
|
||
assert!(output.contains("[输入已关闭]"));
|
||
assert!(output.contains("已完成可信终态"));
|
||
assert!(!output.contains("已退出 Agent Swarm Chat"));
|
||
|
||
fs::remove_dir_all(root).ok();
|
||
}
|
||
|
||
#[test]
|
||
fn recovered_prebaseline_assistant_counts_once_without_duplicate_output() {
|
||
let root = std::env::temp_dir().join(format!(
|
||
"swarm-cli-recovered-assistant-{}-{}",
|
||
std::process::id(),
|
||
unix_millis()
|
||
));
|
||
init_local_game_project_at(
|
||
&root,
|
||
"project-swarm-recovered-assistant",
|
||
"Swarm recovered assistant",
|
||
)
|
||
.expect("initialize recovered assistant project");
|
||
let parent_agent_id = GAME_CREATOR_PROJECT_SUPERVISOR_AGENT_ID;
|
||
let conversation = append_local_conversation_message_at(
|
||
&root,
|
||
Some(parent_agent_id),
|
||
LocalConversationMessage {
|
||
role: "assistant".to_string(),
|
||
content: "已在恢复阶段持久化".to_string(),
|
||
agent_id: Some(parent_agent_id.to_string()),
|
||
},
|
||
)
|
||
.expect("append recovered assistant");
|
||
let session_id = conversation.session_id.expect("active parent session");
|
||
let mut baseline =
|
||
new_swarm_turn_conversation_baseline(conversation.messages.len(), "run-recovered");
|
||
baseline.recovered_assistant = Some(SwarmRecoveredAssistant {
|
||
run_id: "run-recovered".to_string(),
|
||
finalization_id: "finalization-recovered".to_string(),
|
||
message_id: "message-recovered".to_string(),
|
||
content: "已在恢复阶段持久化".to_string(),
|
||
});
|
||
|
||
let snapshot = read_turn_conversation_snapshot(&root, parent_agent_id, &session_id, &baseline)
|
||
.expect("read recovered conversation snapshot");
|
||
assert_eq!(snapshot.metrics.new_assistant_message_count, 1);
|
||
assert_eq!(
|
||
snapshot.metrics.final_reply_chars,
|
||
"已在恢复阶段持久化".chars().count()
|
||
);
|
||
assert_eq!(snapshot.final_reply.as_deref(), Some("已在恢复阶段持久化"));
|
||
assert!(snapshot.recovered_before_observation);
|
||
|
||
let mut output = Vec::new();
|
||
let mut observer = SwarmRuntimeObserver::default();
|
||
let metrics = print_new_parent_reply(
|
||
&root,
|
||
parent_agent_id,
|
||
&session_id,
|
||
&baseline,
|
||
&mut output,
|
||
&mut observer,
|
||
)
|
||
.expect("print recovered parent reply");
|
||
assert_eq!(metrics, snapshot.metrics);
|
||
let output = String::from_utf8(output).expect("recovered output is utf-8");
|
||
assert!(output.contains("父 Agent 回复已在恢复前持久化"));
|
||
assert!(!output.contains("已在恢复阶段持久化"));
|
||
|
||
fs::remove_dir_all(root).ok();
|
||
}
|
||
|
||
#[test]
|
||
fn recovered_assistant_cannot_overlap_a_new_terminal_reply() {
|
||
let root = std::env::temp_dir().join(format!(
|
||
"swarm-cli-recovered-overlap-{}-{}",
|
||
std::process::id(),
|
||
unix_millis()
|
||
));
|
||
init_local_game_project_at(
|
||
&root,
|
||
"project-swarm-recovered-overlap",
|
||
"Swarm recovered overlap",
|
||
)
|
||
.expect("initialize recovered overlap project");
|
||
let parent_agent_id = GAME_CREATOR_PROJECT_SUPERVISOR_AGENT_ID;
|
||
let before =
|
||
read_local_conversation_for_session_at(root.as_path(), Some(parent_agent_id), None)
|
||
.expect("read initial conversation");
|
||
let session_id = before.session_id.expect("active parent session");
|
||
append_local_conversation_message_for_session_at(
|
||
&root,
|
||
Some(parent_agent_id),
|
||
Some(&session_id),
|
||
LocalConversationMessage {
|
||
role: "assistant".to_string(),
|
||
content: "baseline 后的新回复".to_string(),
|
||
agent_id: Some(parent_agent_id.to_string()),
|
||
},
|
||
)
|
||
.expect("append new terminal reply");
|
||
let mut baseline = new_swarm_turn_conversation_baseline(before.messages.len(), "run-overlap");
|
||
baseline.recovered_assistant = Some(SwarmRecoveredAssistant {
|
||
run_id: "run-overlap".to_string(),
|
||
finalization_id: "finalization-overlap".to_string(),
|
||
message_id: "message-overlap".to_string(),
|
||
content: "恢复回复".to_string(),
|
||
});
|
||
|
||
let error = read_turn_conversation_snapshot(&root, parent_agent_id, &session_id, &baseline)
|
||
.expect_err("recovered and new assistant replies must not be double counted");
|
||
assert!(error.contains("与 baseline 后的新回复重叠"));
|
||
|
||
fs::remove_dir_all(root).ok();
|
||
}
|
||
|
||
#[test]
|
||
fn conversation_snapshot_scopes_consecutive_run_replies_by_message_id() {
|
||
let target_message_id = game_creator_agent_runtime_finalization_message_id(
|
||
GAME_CREATOR_PROJECT_SUPERVISOR_AGENT_ID,
|
||
"session-consecutive",
|
||
"run-b",
|
||
);
|
||
let next_message_id = game_creator_agent_runtime_finalization_message_id(
|
||
GAME_CREATOR_PROJECT_SUPERVISOR_AGENT_ID,
|
||
"session-consecutive",
|
||
"run-c",
|
||
);
|
||
let (metrics, reply) = summarize_scoped_new_assistant_messages(
|
||
[
|
||
(
|
||
"assistant",
|
||
"B 的最终回复",
|
||
Some(target_message_id.as_str()),
|
||
),
|
||
("assistant", "C 的最终回复", Some(next_message_id.as_str())),
|
||
],
|
||
Some(&target_message_id),
|
||
);
|
||
assert_eq!(metrics.new_assistant_message_count, 1);
|
||
assert_eq!(metrics.final_reply_chars, "B 的最终回复".chars().count());
|
||
assert_eq!(reply, Some("B 的最终回复"));
|
||
|
||
let (missing_metrics, missing_reply) = summarize_scoped_new_assistant_messages(
|
||
[("assistant", "C 的最终回复", Some(next_message_id.as_str()))],
|
||
Some(&target_message_id),
|
||
);
|
||
assert_eq!(missing_metrics, SwarmTurnConversationMetrics::default());
|
||
assert_eq!(missing_reply, None);
|
||
|
||
let (legacy_metrics, legacy_reply) = summarize_scoped_new_assistant_messages(
|
||
[("assistant", "旧格式回复", None)],
|
||
Some(&target_message_id),
|
||
);
|
||
assert_eq!(legacy_metrics.new_assistant_message_count, 1);
|
||
assert_eq!(legacy_reply, Some("旧格式回复"));
|
||
}
|
||
|
||
#[test]
|
||
fn wait_for_turn_keeps_consecutive_run_reply_and_report_scoped() {
|
||
let root = std::env::temp_dir().join(format!(
|
||
"swarm-cli-consecutive-wait-{}-{}",
|
||
std::process::id(),
|
||
unix_millis()
|
||
));
|
||
init_local_game_project_at(&root, "project-consecutive-wait", "Consecutive wait")
|
||
.expect("initialize consecutive wait project");
|
||
let parent_agent_id = GAME_CREATOR_PROJECT_SUPERVISOR_AGENT_ID;
|
||
let before = append_local_conversation_message_at(
|
||
&root,
|
||
Some(parent_agent_id),
|
||
LocalConversationMessage {
|
||
role: "user".to_string(),
|
||
content: "先完成 B".to_string(),
|
||
agent_id: None,
|
||
},
|
||
)
|
||
.expect("append B user turn");
|
||
let session_id = before.session_id.expect("active supervisor session");
|
||
let baseline = new_swarm_turn_conversation_baseline(before.messages.len(), "run-b");
|
||
for (run_id, reply) in [("run-b", "B 的最终回复"), ("run-c", "C 的最终回复")] {
|
||
append_local_conversation_message_for_session_idempotent_at(
|
||
&root,
|
||
Some(parent_agent_id),
|
||
Some(&session_id),
|
||
LocalConversationMessage {
|
||
role: "assistant".to_string(),
|
||
content: reply.to_string(),
|
||
agent_id: None,
|
||
},
|
||
&game_creator_agent_runtime_finalization_message_id(
|
||
parent_agent_id,
|
||
&session_id,
|
||
run_id,
|
||
),
|
||
)
|
||
.expect("append consecutive assistant reply");
|
||
}
|
||
|
||
let task_path = game_creator_agent_runtime_task_path(&root, parent_agent_id);
|
||
fs::create_dir_all(task_path.parent().expect("consecutive journal parent"))
|
||
.expect("create consecutive journal parent");
|
||
let tasks = [("run-b", "先完成 B"), ("run-c", "再完成 C")]
|
||
.into_iter()
|
||
.map(|(run_id, task)| {
|
||
serde_json::from_value::<AgentRuntimeTaskRecord>(serde_json::json!({
|
||
"agentId": parent_agent_id,
|
||
"taskId": parent_agent_id,
|
||
"sessionId": session_id,
|
||
"runId": run_id,
|
||
"source": AGENT_RUNTIME_SUPERVISOR_CLI_SOURCE,
|
||
"runProfile": AGENT_RUNTIME_RUN_PROFILE_STANDARD,
|
||
"task": task,
|
||
"status": "completed",
|
||
"phase": "completed",
|
||
"currentAction": "本轮已完成",
|
||
"updatedAt": if run_id == "run-b" { 100 } else { 200 }
|
||
}))
|
||
.expect("deserialize consecutive task")
|
||
})
|
||
.collect::<Vec<_>>();
|
||
fs::write(
|
||
&task_path,
|
||
tasks
|
||
.iter()
|
||
.map(|task| serde_json::to_string(task).expect("serialize consecutive task"))
|
||
.collect::<Vec<_>>()
|
||
.join("\n")
|
||
+ "\n",
|
||
)
|
||
.expect("persist consecutive task journal");
|
||
let mut current = default_game_creator_agent_runtime_state(parent_agent_id, "run-c");
|
||
current.session_id = session_id.clone();
|
||
current.source = AGENT_RUNTIME_SUPERVISOR_CLI_SOURCE.to_string();
|
||
current.run_profile = AGENT_RUNTIME_RUN_PROFILE_STANDARD.to_string();
|
||
current.status = "idle".to_string();
|
||
current.phase = "completed".to_string();
|
||
current.current_task = "再完成 C".to_string();
|
||
current.current_action = "C 已完成".to_string();
|
||
current.updated_at = 200;
|
||
write_game_creator_agent_runtime_state(&root, ¤t)
|
||
.expect("persist consecutive canonical state");
|
||
|
||
let (_tx, rx) = mpsc::channel();
|
||
let mut observer = SwarmRuntimeObserver::default();
|
||
let mut output = Vec::new();
|
||
let outcome = wait_for_swarm_turn(
|
||
&root,
|
||
parent_agent_id,
|
||
&session_id,
|
||
AGENT_RUNTIME_RUN_PROFILE_STANDARD,
|
||
None,
|
||
baseline,
|
||
&rx,
|
||
&mut output,
|
||
&mut observer,
|
||
Duration::from_millis(1),
|
||
Duration::ZERO,
|
||
)
|
||
.expect("settle B after canonical advanced to C");
|
||
let SwarmTurnOutcome::Settled(report) = outcome else {
|
||
panic!("B must settle independently: {outcome:?}");
|
||
};
|
||
assert_eq!(report.parent_run_id.as_deref(), Some("run-b"));
|
||
assert_eq!(report.new_assistant_message_count, 1);
|
||
assert_eq!(report.runtime_count, 1);
|
||
let output = String::from_utf8(output).expect("consecutive wait output is utf-8");
|
||
assert!(output.contains("B 的最终回复"));
|
||
assert!(!output.contains("C 的最终回复"));
|
||
|
||
fs::remove_dir_all(root).ok();
|
||
}
|
||
|
||
#[test]
|
||
fn confirmation_prompt_propagates_eof_as_closed_input() {
|
||
let (tx, rx) = mpsc::channel();
|
||
tx.send(SwarmInputEvent::Eof).expect("send eof");
|
||
let mut output = Vec::new();
|
||
|
||
let decision = prompt_swarm_decision(
|
||
Path::new("."),
|
||
GAME_CREATOR_PROJECT_SUPERVISOR_AGENT_ID,
|
||
&rx,
|
||
&mut output,
|
||
"confirm> ",
|
||
)
|
||
.expect("EOF is a turn input state, not an error");
|
||
|
||
assert!(matches!(decision, SwarmPromptDecision::InputClosed));
|
||
}
|
||
|
||
#[test]
|
||
fn terminal_classifier_requires_completed_parent_unique_reply_and_clear_contract() {
|
||
let mut parent = runtime("idle", "completed", 0);
|
||
parent.state.agent_id = GAME_CREATOR_PROJECT_SUPERVISOR_AGENT_ID.to_string();
|
||
parent.state.session_id = "session-terminal".to_string();
|
||
parent.state.run_id = "run-terminal".to_string();
|
||
let unique_reply = SwarmTurnConversationMetrics {
|
||
new_assistant_message_count: 1,
|
||
final_reply_chars: 12,
|
||
};
|
||
|
||
assert_eq!(
|
||
classify_swarm_turn_terminal(Some(&parent), unique_reply, 0, 0, 0),
|
||
SwarmTurnTerminalClassification::Settled
|
||
);
|
||
assert_eq!(
|
||
classify_swarm_turn_terminal(Some(&parent), unique_reply, 0, 1, 0),
|
||
SwarmTurnTerminalClassification::Incomplete
|
||
);
|
||
assert_eq!(
|
||
classify_swarm_turn_terminal(Some(&parent), unique_reply, 0, 0, 1),
|
||
SwarmTurnTerminalClassification::Incomplete
|
||
);
|
||
assert_eq!(
|
||
classify_swarm_turn_terminal(
|
||
Some(&parent),
|
||
SwarmTurnConversationMetrics::default(),
|
||
0,
|
||
0,
|
||
0,
|
||
),
|
||
SwarmTurnTerminalClassification::Incomplete
|
||
);
|
||
assert_eq!(
|
||
classify_swarm_turn_terminal(
|
||
Some(&parent),
|
||
SwarmTurnConversationMetrics {
|
||
new_assistant_message_count: 2,
|
||
final_reply_chars: 12,
|
||
},
|
||
0,
|
||
0,
|
||
0,
|
||
),
|
||
SwarmTurnTerminalClassification::Incomplete
|
||
);
|
||
assert_eq!(
|
||
classify_swarm_turn_terminal(None, unique_reply, 0, 0, 0),
|
||
SwarmTurnTerminalClassification::Incomplete
|
||
);
|
||
}
|
||
|
||
#[test]
|
||
fn terminal_classifier_fails_parent_failure_cancel_and_budget_exhaustion() {
|
||
let metrics = SwarmTurnConversationMetrics {
|
||
new_assistant_message_count: 1,
|
||
final_reply_chars: 8,
|
||
};
|
||
for (status, phase) in [
|
||
("failed", "failed"),
|
||
("cancelled", "cancelled"),
|
||
("failed", "budget-exhausted"),
|
||
] {
|
||
let parent = runtime(status, phase, 0);
|
||
assert_eq!(
|
||
classify_swarm_turn_terminal(Some(&parent), metrics, 0, 0, 0),
|
||
SwarmTurnTerminalClassification::Failed,
|
||
"parent {status}/{phase} must fail closed"
|
||
);
|
||
}
|
||
let completed = runtime("idle", "completed", 0);
|
||
assert_eq!(
|
||
classify_swarm_turn_terminal(Some(&completed), metrics, 1, 0, 0),
|
||
SwarmTurnTerminalClassification::Failed
|
||
);
|
||
}
|
||
|
||
#[test]
|
||
fn pending_interactions_never_form_a_settled_snapshot() {
|
||
let mut parent = runtime("waiting-for-user-input", "waiting-for-user-input", 0);
|
||
parent.state.agent_id = GAME_CREATOR_PROJECT_SUPERVISOR_AGENT_ID.to_string();
|
||
let mut child = runtime("waiting-for-user-input", "waiting-for-user-input", 0);
|
||
child.state.agent_id = "code-prototype".to_string();
|
||
let mut confirmation = runtime("waiting-for-confirmation", "waiting-for-confirmation", 0);
|
||
confirmation.state.agent_id = "quality-review".to_string();
|
||
|
||
assert!(swarm_unhandled_interaction_reasons(
|
||
GAME_CREATOR_PROJECT_SUPERVISOR_AGENT_ID,
|
||
&[parent.clone()],
|
||
false,
|
||
)
|
||
.is_empty());
|
||
let child_reasons = swarm_unhandled_interaction_reasons(
|
||
GAME_CREATOR_PROJECT_SUPERVISOR_AGENT_ID,
|
||
&[child],
|
||
false,
|
||
);
|
||
assert_eq!(child_reasons, vec!["pending-user-input:code-prototype"]);
|
||
let closed_reasons = swarm_unhandled_interaction_reasons(
|
||
GAME_CREATOR_PROJECT_SUPERVISOR_AGENT_ID,
|
||
&[parent, confirmation],
|
||
true,
|
||
);
|
||
assert!(closed_reasons
|
||
.iter()
|
||
.any(|reason| reason == "pending-user-input:project-supervisor"));
|
||
assert!(closed_reasons
|
||
.iter()
|
||
.any(|reason| reason == "pending-confirmation:quality-review"));
|
||
}
|
||
|
||
#[test]
|
||
fn original_specialist_failure_is_recoverable_but_repair_failure_closes() {
|
||
let mut parent = runtime("running", "waiting-for-delegate-receipts", 0);
|
||
parent.state.agent_id = GAME_CREATOR_PROJECT_SUPERVISOR_AGENT_ID.to_string();
|
||
parent.state.session_id = "session-parent".to_string();
|
||
parent.state.run_id = "run-parent".to_string();
|
||
let mut child = runtime("failed", "failed", 0);
|
||
child.state.agent_id = "code-prototype".to_string();
|
||
child.state.session_id = "session-child".to_string();
|
||
child.state.run_id = "run-child".to_string();
|
||
child.state.source = "agent-delegate".to_string();
|
||
child.state.parent_agent_id = Some(parent.state.agent_id.clone());
|
||
child.state.parent_run_id = Some(parent.state.run_id.clone());
|
||
child.state.delegation_id = Some("delivery-original".to_string());
|
||
let acceptance = vec!["交付可运行原型".to_string()];
|
||
let original = new_static_delegate_delivery_with_contract(
|
||
&parent.state.agent_id,
|
||
&parent.state.session_id,
|
||
&parent.state.run_id,
|
||
"action-original",
|
||
"delivery-original",
|
||
&child.state.agent_id,
|
||
&child.state.session_id,
|
||
&child.state.run_id,
|
||
&acceptance,
|
||
&[],
|
||
None,
|
||
);
|
||
|
||
assert_eq!(
|
||
classify_failed_specialist(&parent, &child, Some(&original), false),
|
||
SwarmSpecialistFailureDisposition::Recoverable
|
||
);
|
||
let mut completed_parent = parent.clone();
|
||
completed_parent.state.status = "idle".to_string();
|
||
completed_parent.state.phase = "completed".to_string();
|
||
assert_eq!(
|
||
classify_failed_specialist(&completed_parent, &child, Some(&original), false),
|
||
SwarmSpecialistFailureDisposition::Incomplete
|
||
);
|
||
|
||
child.state.run_id = "run-repair".to_string();
|
||
child.state.delegation_id = Some("delivery-repair".to_string());
|
||
let repair = new_static_delegate_delivery_with_contract(
|
||
&parent.state.agent_id,
|
||
&parent.state.session_id,
|
||
&parent.state.run_id,
|
||
"action-repair",
|
||
"delivery-repair",
|
||
&child.state.agent_id,
|
||
&child.state.session_id,
|
||
&child.state.run_id,
|
||
&acceptance,
|
||
&[],
|
||
Some("delivery-original"),
|
||
);
|
||
assert_eq!(
|
||
classify_failed_specialist(&parent, &child, Some(&repair), false),
|
||
SwarmSpecialistFailureDisposition::Failed
|
||
);
|
||
|
||
let mut successful_repair = repair.clone();
|
||
successful_repair.terminal_status = Some("completed".to_string());
|
||
let mut evidence_ready = StaticDelegateStructuredResult::default();
|
||
evidence_ready.contract_status = StaticDelegateContractStatus::EvidenceReady;
|
||
successful_repair.structured_result = Some(evidence_ready);
|
||
assert!(original_delivery_has_successful_repair(
|
||
&original,
|
||
&[successful_repair.clone()]
|
||
));
|
||
|
||
let mut unknown_original = original.clone();
|
||
let mut unknown_result = StaticDelegateStructuredResult::default();
|
||
unknown_result.contract_status =
|
||
StaticDelegateContractStatus::Unknown("future-contract-status".to_string());
|
||
unknown_original.structured_result = Some(unknown_result);
|
||
assert!(
|
||
!original_delivery_has_successful_repair(&unknown_original, &[successful_repair]),
|
||
"a newer contract status must not be classified as already repaired"
|
||
);
|
||
}
|
||
|
||
#[test]
|
||
fn observer_failure_scan_waits_for_original_repair_and_fails_repair_child() {
|
||
let root = std::env::temp_dir().join(format!(
|
||
"swarm-cli-repair-scan-{}-{}",
|
||
std::process::id(),
|
||
unix_millis()
|
||
));
|
||
init_local_game_project_at(&root, "project-swarm-repair", "Swarm repair scan")
|
||
.expect("initialize repair scan project");
|
||
let mut parent = runtime("running", "waiting-for-delegate-receipts", 0);
|
||
parent.state.agent_id = GAME_CREATOR_PROJECT_SUPERVISOR_AGENT_ID.to_string();
|
||
parent.state.session_id = "session-parent".to_string();
|
||
parent.state.run_id = "run-parent".to_string();
|
||
let mut child = runtime("failed", "failed", 0);
|
||
child.state.agent_id = "code-prototype".to_string();
|
||
child.state.session_id = "session-child".to_string();
|
||
child.state.run_id = "run-child".to_string();
|
||
child.state.source = "agent-delegate".to_string();
|
||
child.state.parent_agent_id = Some(parent.state.agent_id.clone());
|
||
child.state.parent_run_id = Some(parent.state.run_id.clone());
|
||
child.state.delegation_id = Some("delivery-original".to_string());
|
||
let acceptance = vec!["交付可运行原型".to_string()];
|
||
let original = new_static_delegate_delivery_with_contract(
|
||
&parent.state.agent_id,
|
||
&parent.state.session_id,
|
||
&parent.state.run_id,
|
||
"action-original",
|
||
"delivery-original",
|
||
&child.state.agent_id,
|
||
&child.state.session_id,
|
||
&child.state.run_id,
|
||
&acceptance,
|
||
&[],
|
||
None,
|
||
);
|
||
create_or_read_static_delegate_delivery_at(&root, &original)
|
||
.expect("persist original delivery");
|
||
|
||
let original_scan = scan_swarm_terminal_failures_at(
|
||
&root,
|
||
&parent.state.agent_id,
|
||
&parent.state.session_id,
|
||
AGENT_RUNTIME_RUN_PROFILE_STANDARD,
|
||
None,
|
||
&parent.state.run_id,
|
||
&[parent.clone(), child.clone()],
|
||
);
|
||
assert!(original_scan.failed_agents.is_empty());
|
||
assert!(original_scan.incomplete_reasons.is_empty());
|
||
assert!(original_scan.reconciliation_agents.is_empty());
|
||
|
||
child.state.run_id = "run-repair".to_string();
|
||
child.state.delegation_id = Some("delivery-repair".to_string());
|
||
let repair = new_static_delegate_delivery_with_contract(
|
||
&parent.state.agent_id,
|
||
&parent.state.session_id,
|
||
&parent.state.run_id,
|
||
"action-repair",
|
||
"delivery-repair",
|
||
&child.state.agent_id,
|
||
&child.state.session_id,
|
||
&child.state.run_id,
|
||
&acceptance,
|
||
&[],
|
||
Some("delivery-original"),
|
||
);
|
||
create_or_read_static_delegate_delivery_at(&root, &repair).expect("persist repair delivery");
|
||
let repair_scan = scan_swarm_terminal_failures_at(
|
||
&root,
|
||
&parent.state.agent_id,
|
||
&parent.state.session_id,
|
||
AGENT_RUNTIME_RUN_PROFILE_STANDARD,
|
||
None,
|
||
&parent.state.run_id,
|
||
&[parent.clone(), child],
|
||
);
|
||
assert_eq!(repair_scan.failed_agents, vec!["code-prototype:failed"]);
|
||
assert!(repair_scan.incomplete_reasons.is_empty());
|
||
assert!(repair_scan.reconciliation_agents.is_empty());
|
||
|
||
fs::remove_dir_all(root).ok();
|
||
}
|
||
|
||
#[test]
|
||
fn failure_scan_reads_all_historical_runs_for_the_same_specialist() {
|
||
let root = std::env::temp_dir().join(format!(
|
||
"swarm-cli-historical-repair-scan-{}-{}",
|
||
std::process::id(),
|
||
unix_millis()
|
||
));
|
||
init_local_game_project_at(&root, "project-historical-repair", "Historical repair scan")
|
||
.expect("initialize historical repair project");
|
||
let mut parent = runtime("running", "waiting-for-delegate-receipts", 0);
|
||
parent.state.agent_id = GAME_CREATOR_PROJECT_SUPERVISOR_AGENT_ID.to_string();
|
||
parent.state.session_id = "session-historical-parent".to_string();
|
||
parent.state.run_id = "run-historical-parent".to_string();
|
||
let acceptance = vec!["交付可运行原型".to_string()];
|
||
let original = new_static_delegate_delivery_with_contract(
|
||
&parent.state.agent_id,
|
||
&parent.state.session_id,
|
||
&parent.state.run_id,
|
||
"action-historical-original",
|
||
"delivery-historical-original",
|
||
"code-prototype",
|
||
"session-historical-child",
|
||
"run-historical-original",
|
||
&acceptance,
|
||
&[],
|
||
None,
|
||
);
|
||
let repair = new_static_delegate_delivery_with_contract(
|
||
&parent.state.agent_id,
|
||
&parent.state.session_id,
|
||
&parent.state.run_id,
|
||
"action-historical-repair",
|
||
"delivery-historical-repair",
|
||
"code-prototype",
|
||
"session-historical-child",
|
||
"run-historical-repair",
|
||
&acceptance,
|
||
&[],
|
||
Some("delivery-historical-original"),
|
||
);
|
||
let quality_failure = new_static_delegate_delivery_with_contract(
|
||
&parent.state.agent_id,
|
||
&parent.state.session_id,
|
||
&parent.state.run_id,
|
||
"action-historical-quality",
|
||
"delivery-historical-quality",
|
||
"quality-review",
|
||
"session-historical-quality",
|
||
"run-historical-repair",
|
||
&acceptance,
|
||
&[],
|
||
Some("delivery-historical-quality-original"),
|
||
);
|
||
create_or_read_static_delegate_delivery_at(&root, &original)
|
||
.expect("persist historical original delivery");
|
||
create_or_read_static_delegate_delivery_at(&root, &repair)
|
||
.expect("persist historical repair delivery");
|
||
create_or_read_static_delegate_delivery_at(&root, &quality_failure)
|
||
.expect("persist historical quality delivery");
|
||
|
||
let original_task = serde_json::from_value::<AgentRuntimeTaskRecord>(serde_json::json!({
|
||
"agentId": "code-prototype",
|
||
"taskId": "code-prototype",
|
||
"sessionId": "session-historical-child",
|
||
"runId": "run-historical-original",
|
||
"source": "agent-delegate",
|
||
"runProfile": AGENT_RUNTIME_RUN_PROFILE_STANDARD,
|
||
"parentAgentId": GAME_CREATOR_PROJECT_SUPERVISOR_AGENT_ID,
|
||
"parentRunId": "run-historical-parent",
|
||
"delegationId": "delivery-historical-original",
|
||
"task": "原始委派",
|
||
"status": "failed",
|
||
"phase": "failed",
|
||
"currentAction": "原始委派失败",
|
||
"updatedAt": 100
|
||
}))
|
||
.expect("deserialize historical original task");
|
||
let repair_task = serde_json::from_value::<AgentRuntimeTaskRecord>(serde_json::json!({
|
||
"agentId": "code-prototype",
|
||
"taskId": "code-prototype",
|
||
"sessionId": "session-historical-child",
|
||
"runId": "run-historical-repair",
|
||
"source": "agent-delegate",
|
||
"runProfile": AGENT_RUNTIME_RUN_PROFILE_STANDARD,
|
||
"parentAgentId": GAME_CREATOR_PROJECT_SUPERVISOR_AGENT_ID,
|
||
"parentRunId": "run-historical-parent",
|
||
"delegationId": "delivery-historical-repair",
|
||
"task": "修复委派",
|
||
"status": "failed",
|
||
"phase": "completion-contract-failed",
|
||
"currentAction": "修复交付未通过合同",
|
||
"updatedAt": 200
|
||
}))
|
||
.expect("deserialize historical repair task");
|
||
let task_path = game_creator_agent_runtime_task_path(&root, "code-prototype");
|
||
fs::create_dir_all(task_path.parent().expect("specialist journal parent"))
|
||
.expect("create specialist journal parent");
|
||
fs::write(
|
||
&task_path,
|
||
format!(
|
||
"{}\n{}\n",
|
||
serde_json::to_string(&original_task).expect("serialize original task"),
|
||
serde_json::to_string(&repair_task).expect("serialize repair task"),
|
||
),
|
||
)
|
||
.expect("persist specialist task journal");
|
||
|
||
let quality_task = serde_json::from_value::<AgentRuntimeTaskRecord>(serde_json::json!({
|
||
"agentId": "quality-review",
|
||
"taskId": "quality-review",
|
||
"sessionId": "session-historical-quality",
|
||
"runId": "run-historical-repair",
|
||
"source": "agent-delegate",
|
||
"runProfile": AGENT_RUNTIME_RUN_PROFILE_STANDARD,
|
||
"parentAgentId": GAME_CREATOR_PROJECT_SUPERVISOR_AGENT_ID,
|
||
"parentRunId": "run-historical-parent",
|
||
"delegationId": "delivery-historical-quality",
|
||
"task": "质量修复复验",
|
||
"status": "failed",
|
||
"phase": "failed",
|
||
"currentAction": "质量修复失败",
|
||
"updatedAt": 210
|
||
}))
|
||
.expect("deserialize historical quality task");
|
||
let quality_task_path = game_creator_agent_runtime_task_path(&root, "quality-review");
|
||
fs::create_dir_all(quality_task_path.parent().expect("quality journal parent"))
|
||
.expect("create quality journal parent");
|
||
fs::write(
|
||
&quality_task_path,
|
||
format!(
|
||
"{}\n",
|
||
serde_json::to_string(&quality_task).expect("serialize quality task"),
|
||
),
|
||
)
|
||
.expect("persist quality task journal");
|
||
|
||
let mut current_specialist = runtime("running", "planning", 0);
|
||
current_specialist.state.agent_id = "code-prototype".to_string();
|
||
current_specialist.state.session_id = "session-historical-child".to_string();
|
||
current_specialist.state.run_id = "run-historical-repair".to_string();
|
||
current_specialist.state.source = "agent-delegate".to_string();
|
||
current_specialist.state.parent_agent_id = Some(parent.state.agent_id.clone());
|
||
current_specialist.state.parent_run_id = Some(parent.state.run_id.clone());
|
||
current_specialist.state.updated_at = 50;
|
||
let mut current_quality = runtime("idle", "completed", 0);
|
||
current_quality.state.agent_id = "quality-review".to_string();
|
||
current_quality.state.run_id = "run-later-quality".to_string();
|
||
let scan = scan_swarm_terminal_failures_at(
|
||
&root,
|
||
&parent.state.agent_id,
|
||
&parent.state.session_id,
|
||
AGENT_RUNTIME_RUN_PROFILE_STANDARD,
|
||
None,
|
||
&parent.state.run_id,
|
||
&[parent.clone(), current_specialist, current_quality],
|
||
);
|
||
assert_eq!(
|
||
scan.failed_agents,
|
||
vec![
|
||
"code-prototype:completion-contract-failed",
|
||
"quality-review:failed",
|
||
]
|
||
);
|
||
assert!(scan.incomplete_reasons.is_empty());
|
||
assert!(scan.reconciliation_agents.is_empty());
|
||
|
||
fs::remove_dir_all(root).ok();
|
||
}
|
||
|
||
#[test]
|
||
fn failure_scan_ignores_cancelled_parent_and_children_from_another_run() {
|
||
let root = std::env::temp_dir().join(format!(
|
||
"swarm-cli-cross-run-failure-scan-{}-{}",
|
||
std::process::id(),
|
||
unix_millis()
|
||
));
|
||
init_local_game_project_at(&root, "project-cross-run-scan", "Cross-run failure scan")
|
||
.expect("initialize cross-run failure scan project");
|
||
let session_id = "session-cross-run";
|
||
let run_profile = AGENT_RUNTIME_RUN_PROFILE_AUTONOMOUS_GAME_BUILD;
|
||
let mut old_parent = runtime("cancelled", "cancelled", 1);
|
||
old_parent.state.agent_id = GAME_CREATOR_PROJECT_SUPERVISOR_AGENT_ID.to_string();
|
||
old_parent.state.session_id = session_id.to_string();
|
||
old_parent.state.run_id = "run-old-cancelled".to_string();
|
||
old_parent.state.run_profile = run_profile.to_string();
|
||
|
||
let mut old_child = runtime("failed", "budget-exhausted", 0);
|
||
old_child.state.agent_id = "art-asset-plan".to_string();
|
||
old_child.state.session_id = "session-old-child".to_string();
|
||
old_child.state.run_id = "run-old-child".to_string();
|
||
old_child.state.source = "agent-delegate".to_string();
|
||
old_child.state.parent_agent_id = Some(old_parent.state.agent_id.clone());
|
||
old_child.state.parent_run_id = Some(old_parent.state.run_id.clone());
|
||
|
||
let runtimes = vec![old_parent.clone(), old_child];
|
||
let new_turn_scan = scan_swarm_terminal_failures_at(
|
||
&root,
|
||
&old_parent.state.agent_id,
|
||
session_id,
|
||
run_profile,
|
||
None,
|
||
"run-new-pending",
|
||
&runtimes,
|
||
);
|
||
assert!(new_turn_scan.failed_agents.is_empty());
|
||
assert!(new_turn_scan.incomplete_reasons.is_empty());
|
||
assert!(new_turn_scan.reconciliation_agents.is_empty());
|
||
|
||
let old_turn_scan = scan_swarm_terminal_failures_at(
|
||
&root,
|
||
&old_parent.state.agent_id,
|
||
session_id,
|
||
run_profile,
|
||
None,
|
||
&old_parent.state.run_id,
|
||
&runtimes,
|
||
);
|
||
assert!(old_turn_scan
|
||
.failed_agents
|
||
.iter()
|
||
.any(|agent| agent == "project-supervisor:cancelled"));
|
||
assert!(old_turn_scan
|
||
.failed_agents
|
||
.iter()
|
||
.any(|agent| agent == "art-asset-plan:budget-exhausted"));
|
||
|
||
fs::remove_dir_all(root).ok();
|
||
}
|
||
|
||
#[test]
|
||
fn missing_confirmation_sidecar_is_reported_as_reconciliation() {
|
||
let broken = runtime("waiting-for-confirmation", "waiting-for-confirmation", 0);
|
||
assert_eq!(
|
||
swarm_reconciliation_agents(&[broken]),
|
||
vec!["code-prototype".to_string()]
|
||
);
|
||
}
|
||
|
||
#[test]
|
||
fn turn_report_counts_runtime_and_conversation_snapshots() {
|
||
let mut parent = runtime("running", "response", 0);
|
||
parent.state.agent_id = "project-supervisor".to_string();
|
||
parent.state.session_id = "session-report".to_string();
|
||
parent.state.run_id = "run-parent".to_string();
|
||
|
||
let mut pending_child = runtime("pending", "queued", 0);
|
||
pending_child.state.agent_id = "child-code".to_string();
|
||
pending_child.state.run_id = "run-child-pending".to_string();
|
||
pending_child.state.parent_agent_id = Some("project-supervisor".to_string());
|
||
pending_child.state.parent_run_id = Some("run-parent".to_string());
|
||
|
||
let mut confirmation_child = runtime("waiting-for-confirmation", "waiting-for-confirmation", 0);
|
||
confirmation_child.state.agent_id = "child-design".to_string();
|
||
confirmation_child.state.run_id = "run-child-confirmation".to_string();
|
||
confirmation_child.state.parent_agent_id = Some("project-supervisor".to_string());
|
||
confirmation_child.state.parent_run_id = Some("run-parent".to_string());
|
||
|
||
let mut input_child = runtime("waiting-for-user-input", "waiting-for-user-input", 0);
|
||
input_child.state.agent_id = "child-test".to_string();
|
||
input_child.state.run_id = "run-child-input".to_string();
|
||
input_child.state.parent_agent_id = Some("project-supervisor".to_string());
|
||
input_child.state.parent_run_id = Some("run-parent".to_string());
|
||
|
||
let mut next_turn = runtime("running", "planning", 0);
|
||
next_turn.state.agent_id = "project-supervisor".to_string();
|
||
next_turn.state.session_id = "session-report".to_string();
|
||
next_turn.state.run_id = "run-next-turn".to_string();
|
||
let runtimes = vec![
|
||
parent,
|
||
pending_child,
|
||
confirmation_child,
|
||
input_child,
|
||
next_turn,
|
||
];
|
||
let (conversation_metrics, final_reply) = summarize_new_assistant_messages([
|
||
("user", "请继续"),
|
||
("assistant", "阶段回复"),
|
||
("tool", "PRIVATE_OBSERVATION"),
|
||
("assistant", "最终🙂"),
|
||
]);
|
||
assert_eq!(final_reply, Some("最终🙂"));
|
||
|
||
let report = build_swarm_turn_report(
|
||
SwarmTurnReportOutcome::NeedsReconciliation,
|
||
"project-supervisor",
|
||
"session-report",
|
||
Some("run-parent"),
|
||
&runtimes,
|
||
conversation_metrics,
|
||
1,
|
||
)
|
||
.expect("build scoped turn report");
|
||
|
||
assert_eq!(report.schema_version, SWARM_TURN_REPORT_SCHEMA_VERSION);
|
||
assert_eq!(report.outcome, SwarmTurnReportOutcome::NeedsReconciliation);
|
||
assert_eq!(report.parent_agent_id, "project-supervisor");
|
||
assert_eq!(report.session_id, "session-report");
|
||
assert_eq!(report.parent_run_id.as_deref(), Some("run-parent"));
|
||
assert_eq!(report.runtime_count, 4);
|
||
assert_eq!(report.busy_runtime_count, 4);
|
||
assert_eq!(report.pending_task_count, 1);
|
||
assert_eq!(report.running_task_count, 1);
|
||
assert_eq!(report.waiting_for_confirmation_count, 1);
|
||
assert_eq!(report.waiting_for_user_input_count, 1);
|
||
assert_eq!(report.new_assistant_message_count, 2);
|
||
assert_eq!(report.final_reply_chars, "最终🙂".chars().count());
|
||
assert_eq!(report.reconciliation_agent_count, 1);
|
||
}
|
||
|
||
#[test]
|
||
fn turn_report_prefers_expected_run_over_stale_canonical_state() {
|
||
let root = std::env::temp_dir().join(format!(
|
||
"swarm-cli-report-full-journal-{}-{}",
|
||
std::process::id(),
|
||
unix_millis()
|
||
));
|
||
init_local_game_project_at(&root, "project-report-journal", "Report full journal")
|
||
.expect("initialize report journal project");
|
||
let mut old_parent = runtime("cancelled", "cancelled", 1);
|
||
old_parent.state.agent_id = GAME_CREATOR_PROJECT_SUPERVISOR_AGENT_ID.to_string();
|
||
old_parent.state.session_id = "session-report-stale".to_string();
|
||
old_parent.state.run_id = "run-old-cancelled".to_string();
|
||
let task_path =
|
||
game_creator_agent_runtime_task_path(&root, GAME_CREATOR_PROJECT_SUPERVISOR_AGENT_ID);
|
||
fs::create_dir_all(task_path.parent().expect("report journal parent"))
|
||
.expect("create report journal parent");
|
||
let pending_task = serde_json::from_value::<AgentRuntimeTaskRecord>(serde_json::json!({
|
||
"agentId": GAME_CREATOR_PROJECT_SUPERVISOR_AGENT_ID,
|
||
"taskId": GAME_CREATOR_PROJECT_SUPERVISOR_AGENT_ID,
|
||
"sessionId": "session-report-stale",
|
||
"runId": "run-new-pending",
|
||
"source": AGENT_RUNTIME_SUPERVISOR_CLI_SOURCE,
|
||
"runProfile": AGENT_RUNTIME_RUN_PROFILE_STANDARD,
|
||
"task": "等待下一轮",
|
||
"status": "pending",
|
||
"phase": "queued",
|
||
"currentAction": "等待 Runner",
|
||
"updatedAt": 200
|
||
}))
|
||
.expect("deserialize report pending task");
|
||
fs::write(
|
||
&task_path,
|
||
format!(
|
||
"{}\n",
|
||
serde_json::to_string(&pending_task).expect("serialize report pending task"),
|
||
),
|
||
)
|
||
.expect("persist report task journal");
|
||
old_parent.task_path = task_path.to_string_lossy().into_owned();
|
||
|
||
let report = build_swarm_turn_report(
|
||
SwarmTurnReportOutcome::Failed,
|
||
GAME_CREATOR_PROJECT_SUPERVISOR_AGENT_ID,
|
||
"session-report-stale",
|
||
Some("run-new-pending"),
|
||
&[old_parent],
|
||
SwarmTurnConversationMetrics::default(),
|
||
0,
|
||
)
|
||
.expect("build stale canonical turn report");
|
||
|
||
assert_eq!(report.parent_run_id.as_deref(), Some("run-new-pending"));
|
||
assert_eq!(report.runtime_count, 1);
|
||
assert_eq!(report.pending_task_count, 1);
|
||
|
||
fs::remove_dir_all(root).ok();
|
||
}
|
||
|
||
#[test]
|
||
fn turn_report_json_is_single_line_and_omits_sensitive_bodies_and_paths() {
|
||
let sensitive_reply = concat!(
|
||
"PRIVATE_REPLY_BODY\n",
|
||
"/private/project/root ",
|
||
"prompt=DO_NOT_LEAK observation=DO_NOT_LEAK CREDENTIAL_SENTINEL"
|
||
);
|
||
let (conversation_metrics, _) =
|
||
summarize_new_assistant_messages([("assistant", sensitive_reply)]);
|
||
let report = build_swarm_turn_report(
|
||
SwarmTurnReportOutcome::Settled,
|
||
"project-supervisor",
|
||
"session-safe",
|
||
None,
|
||
&[],
|
||
conversation_metrics,
|
||
0,
|
||
)
|
||
.expect("build safe turn report");
|
||
let json = serde_json::to_string(&report).expect("serialize turn report");
|
||
let value = serde_json::from_str::<serde_json::Value>(&json).expect("parse turn report");
|
||
let object = value.as_object().expect("turn report is an object");
|
||
|
||
assert_eq!(json.lines().count(), 1);
|
||
assert_eq!(object.len(), 14);
|
||
for key in [
|
||
"schemaVersion",
|
||
"outcome",
|
||
"parentAgentId",
|
||
"sessionId",
|
||
"parentRunId",
|
||
"runtimeCount",
|
||
"busyRuntimeCount",
|
||
"pendingTaskCount",
|
||
"runningTaskCount",
|
||
"waitingForConfirmationCount",
|
||
"waitingForUserInputCount",
|
||
"newAssistantMessageCount",
|
||
"finalReplyChars",
|
||
"reconciliationAgentCount",
|
||
] {
|
||
assert!(object.contains_key(key), "turn report omitted {key}");
|
||
}
|
||
assert_eq!(
|
||
value["schemaVersion"],
|
||
serde_json::json!(SWARM_TURN_REPORT_SCHEMA_VERSION)
|
||
);
|
||
assert_eq!(value["outcome"], serde_json::json!("settled"));
|
||
assert_eq!(value["parentRunId"], serde_json::Value::Null);
|
||
assert_eq!(value["newAssistantMessageCount"], serde_json::json!(1));
|
||
assert_eq!(
|
||
value["finalReplyChars"],
|
||
serde_json::json!(sensitive_reply.chars().count())
|
||
);
|
||
for forbidden in [
|
||
"PRIVATE_REPLY_BODY",
|
||
"/private/project/root",
|
||
"DO_NOT_LEAK",
|
||
"CREDENTIAL_SENTINEL",
|
||
] {
|
||
assert!(!json.contains(forbidden), "report leaked {forbidden}");
|
||
}
|
||
}
|
||
|
||
#[test]
|
||
fn turn_outcome_prints_all_terminal_reports_but_not_quit() {
|
||
let metrics = SwarmTurnConversationMetrics {
|
||
new_assistant_message_count: 1,
|
||
final_reply_chars: 4,
|
||
};
|
||
let settled_report = build_swarm_turn_report(
|
||
SwarmTurnReportOutcome::Settled,
|
||
"project-supervisor",
|
||
"session-settled",
|
||
None,
|
||
&[],
|
||
metrics,
|
||
0,
|
||
)
|
||
.expect("build settled turn report");
|
||
let mut settled_output = Vec::new();
|
||
print_turn_outcome(
|
||
SwarmTurnOutcome::Settled(settled_report),
|
||
&mut settled_output,
|
||
)
|
||
.expect("print settled report");
|
||
let settled_output = String::from_utf8(settled_output).expect("settled output is utf-8");
|
||
assert_eq!(settled_output.lines().count(), 1);
|
||
assert!(settled_output.starts_with(SWARM_TURN_REPORT_PREFIX));
|
||
assert!(settled_output.contains("\"outcome\":\"settled\""));
|
||
|
||
let failed_report = build_swarm_turn_report(
|
||
SwarmTurnReportOutcome::Failed,
|
||
"project-supervisor",
|
||
"session-failed",
|
||
None,
|
||
&[],
|
||
metrics,
|
||
0,
|
||
)
|
||
.expect("build failed turn report");
|
||
let mut failed_output = Vec::new();
|
||
print_turn_outcome(
|
||
SwarmTurnOutcome::Failed {
|
||
agent_ids: vec!["project-supervisor:budget-exhausted".to_string()],
|
||
report: failed_report,
|
||
},
|
||
&mut failed_output,
|
||
)
|
||
.expect("print failed report");
|
||
let failed_output = String::from_utf8(failed_output).expect("failed output is utf-8");
|
||
assert!(failed_output.starts_with("[已失败]"));
|
||
assert!(failed_output.contains("\"outcome\":\"failed\""));
|
||
|
||
let incomplete_report = build_swarm_turn_report(
|
||
SwarmTurnReportOutcome::Incomplete,
|
||
"project-supervisor",
|
||
"session-incomplete",
|
||
None,
|
||
&[],
|
||
metrics,
|
||
0,
|
||
)
|
||
.expect("build incomplete turn report");
|
||
let mut incomplete_output = Vec::new();
|
||
print_turn_outcome(
|
||
SwarmTurnOutcome::Incomplete {
|
||
reasons: vec!["assistant-count=0".to_string()],
|
||
report: incomplete_report,
|
||
},
|
||
&mut incomplete_output,
|
||
)
|
||
.expect("print incomplete report");
|
||
let incomplete_output =
|
||
String::from_utf8(incomplete_output).expect("incomplete output is utf-8");
|
||
assert!(incomplete_output.starts_with("[未完成]"));
|
||
assert!(incomplete_output.contains("\"outcome\":\"incomplete\""));
|
||
|
||
let reconciliation_report = build_swarm_turn_report(
|
||
SwarmTurnReportOutcome::NeedsReconciliation,
|
||
"project-supervisor",
|
||
"session-reconciliation",
|
||
None,
|
||
&[],
|
||
metrics,
|
||
2,
|
||
)
|
||
.expect("build reconciliation turn report");
|
||
let mut reconciliation_output = Vec::new();
|
||
print_turn_outcome(
|
||
SwarmTurnOutcome::NeedsReconciliation {
|
||
agent_ids: vec!["code-prototype".to_string(), "external-runner".to_string()],
|
||
report: reconciliation_report,
|
||
},
|
||
&mut reconciliation_output,
|
||
)
|
||
.expect("print reconciliation report");
|
||
let reconciliation_output =
|
||
String::from_utf8(reconciliation_output).expect("reconciliation output is utf-8");
|
||
let lines = reconciliation_output.lines().collect::<Vec<_>>();
|
||
assert_eq!(lines.len(), 2);
|
||
assert_eq!(
|
||
lines[0],
|
||
"[已阻断] 以下 Agent 需要人工 reconciliation:code-prototype, external-runner"
|
||
);
|
||
assert!(lines[1].starts_with(SWARM_TURN_REPORT_PREFIX));
|
||
assert!(lines[1].contains("\"outcome\":\"needs-reconciliation\""));
|
||
assert!(lines[1].contains("\"reconciliationAgentCount\":2"));
|
||
|
||
let mut quit_output = Vec::new();
|
||
print_turn_outcome(SwarmTurnOutcome::Quit, &mut quit_output).expect("ignore quit");
|
||
assert!(quit_output.is_empty());
|
||
}
|
||
|
||
#[test]
|
||
fn blank_agent_snapshots_do_not_reset_the_settle_window() {
|
||
let mut blank = runtime("idle", "idle", 0);
|
||
blank.state.run_id.clear();
|
||
blank.state.updated_at = 100;
|
||
let mut observer = SwarmRuntimeObserver::default();
|
||
let mut output = Vec::new();
|
||
assert!(!observer
|
||
.print_changes(&[blank.clone()], &mut output)
|
||
.expect("observe first blank snapshot"));
|
||
blank.state.updated_at = 101;
|
||
assert!(!observer
|
||
.print_changes(&[blank], &mut output)
|
||
.expect("observe refreshed blank snapshot"));
|
||
assert!(output.is_empty());
|
||
}
|
||
|
||
#[test]
|
||
fn runtime_plan_revision_and_current_step_change_state_signature() {
|
||
let mut snapshot = runtime("running", "planning", 0);
|
||
snapshot.state.updated_at = 100;
|
||
snapshot.task_queue.updated_at = 100;
|
||
snapshot.state.plan_revision = 1;
|
||
snapshot.state.plan_steps = vec![
|
||
AgentRuntimePlanStep {
|
||
index: 0,
|
||
title: "读取现有 CLI".to_string(),
|
||
status: "in_progress".to_string(),
|
||
detail: None,
|
||
updated_at: 100,
|
||
},
|
||
AgentRuntimePlanStep {
|
||
index: 1,
|
||
title: "补充计划展示".to_string(),
|
||
status: "pending".to_string(),
|
||
detail: None,
|
||
updated_at: 100,
|
||
},
|
||
];
|
||
snapshot.state.active_plan_step_index = Some(0);
|
||
|
||
let initial = runtime_state_signature(&snapshot.state, &snapshot.task_queue);
|
||
snapshot.state.plan_revision = 2;
|
||
let revised = runtime_state_signature(&snapshot.state, &snapshot.task_queue);
|
||
assert_ne!(initial, revised);
|
||
|
||
snapshot.state.plan_steps[0].title = "核对现有 CLI".to_string();
|
||
let current_step_changed = runtime_state_signature(&snapshot.state, &snapshot.task_queue);
|
||
assert_ne!(revised, current_step_changed);
|
||
|
||
snapshot.state.plan_steps[0].status = "completed".to_string();
|
||
snapshot.state.plan_steps[1].status = "in_progress".to_string();
|
||
snapshot.state.active_plan_step_index = Some(1);
|
||
let advanced = runtime_state_signature(&snapshot.state, &snapshot.task_queue);
|
||
assert_ne!(current_step_changed, advanced);
|
||
}
|
||
|
||
#[test]
|
||
fn runtime_plan_output_is_bounded_and_omits_private_observations() {
|
||
let mut snapshot = runtime("running", "planning", 0);
|
||
snapshot.state.plan_revision = 7;
|
||
snapshot.state.plan_explanation = "已完成读取,进入验证".to_string();
|
||
snapshot.state.current_action = "展示持久计划".to_string();
|
||
snapshot.state.waiting_on = "开发者确认".to_string();
|
||
snapshot.state.next_step = "运行 focused cargo test".to_string();
|
||
snapshot.state.observations = vec![
|
||
"PRIVATE_OBSERVATION_SENTINEL".to_string(),
|
||
"PRIVATE_DETAIL_SENTINEL".to_string(),
|
||
];
|
||
snapshot.state.plan_steps = (0..10)
|
||
.map(|index| AgentRuntimePlanStep {
|
||
index,
|
||
title: format!("计划步骤 {}", index + 1),
|
||
status: match index {
|
||
0 | 1 => "completed",
|
||
2 => "in_progress",
|
||
_ => "pending",
|
||
}
|
||
.to_string(),
|
||
detail: Some(format!("PRIVATE_STEP_DETAIL_{index}")),
|
||
updated_at: 100,
|
||
})
|
||
.collect();
|
||
snapshot.state.active_plan_step_index = Some(2);
|
||
|
||
let mut output = Vec::new();
|
||
print_runtime_state(&snapshot.state, &snapshot.task_queue, &mut output)
|
||
.expect("print runtime plan progress");
|
||
let output = String::from_utf8(output).expect("runtime output is utf-8");
|
||
|
||
assert!(output.contains(
|
||
"[计划] revision=7 completed=2/10 current=#3 [in_progress] 计划步骤 3 | waiting=开发者确认 | next=运行 focused cargo test"
|
||
));
|
||
assert!(output.contains("[计划说明] 已完成读取,进入验证"));
|
||
assert_eq!(output.matches("[计划步骤]").count(), 8);
|
||
assert!(output.contains("[计划步骤] #8 [pending] 计划步骤 8"));
|
||
assert!(output.contains("另有 2 条步骤未显示"));
|
||
assert!(!output.contains("计划步骤 9"));
|
||
assert!(!output.contains("PRIVATE_OBSERVATION_SENTINEL"));
|
||
assert!(!output.contains("PRIVATE_DETAIL_SENTINEL"));
|
||
assert!(!output.contains("PRIVATE_STEP_DETAIL"));
|
||
}
|
||
|
||
#[test]
|
||
fn response_stream_prints_only_monotonic_utf8_suffixes() {
|
||
let mut observer = SwarmRuntimeObserver::default();
|
||
let mut output = Vec::new();
|
||
let mut snapshot = runtime_with_response_stream(response_stream(
|
||
"slot-1",
|
||
7,
|
||
0,
|
||
AGENT_RUNTIME_RESPONSE_STREAM_STATUS_STREAMING,
|
||
"",
|
||
));
|
||
|
||
assert!(observer
|
||
.print_changes(&[snapshot.clone()], &mut output)
|
||
.expect("observe empty response stream"));
|
||
snapshot.response_stream = Some(response_stream(
|
||
"slot-1",
|
||
7,
|
||
1,
|
||
AGENT_RUNTIME_RESPONSE_STREAM_STATUS_STREAMING,
|
||
"你",
|
||
));
|
||
assert!(observer
|
||
.print_changes(&[snapshot.clone()], &mut output)
|
||
.expect("observe first utf-8 suffix"));
|
||
snapshot.response_stream = Some(response_stream(
|
||
"slot-1",
|
||
7,
|
||
2,
|
||
AGENT_RUNTIME_RESPONSE_STREAM_STATUS_READY,
|
||
"你好🙂",
|
||
));
|
||
assert!(observer
|
||
.print_changes(&[snapshot.clone()], &mut output)
|
||
.expect("observe second utf-8 suffix"));
|
||
assert!(!observer
|
||
.print_changes(&[snapshot], &mut output)
|
||
.expect("ignore duplicate snapshot"));
|
||
observer
|
||
.close_response_line(&mut output)
|
||
.expect("close response line");
|
||
|
||
let output = String::from_utf8(output).expect("stream output is utf-8");
|
||
assert!(output.contains("Agent[code-prototype]> 你好🙂"));
|
||
assert_eq!(output.matches("Agent[code-prototype]>").count(), 1);
|
||
assert!(!output.contains("你你好"));
|
||
}
|
||
|
||
#[test]
|
||
fn response_stream_resets_for_non_prefix_and_new_request_slot() {
|
||
let mut observer = SwarmRuntimeObserver::default();
|
||
let mut output = Vec::new();
|
||
let mut snapshot = runtime_with_response_stream(response_stream(
|
||
"slot-1",
|
||
9,
|
||
1,
|
||
AGENT_RUNTIME_RESPONSE_STREAM_STATUS_STREAMING,
|
||
"旧稿",
|
||
));
|
||
observer
|
||
.print_changes(&[snapshot.clone()], &mut output)
|
||
.expect("observe initial stream");
|
||
|
||
snapshot.response_stream = Some(response_stream(
|
||
"slot-1",
|
||
9,
|
||
2,
|
||
AGENT_RUNTIME_RESPONSE_STREAM_STATUS_READY,
|
||
"修正版",
|
||
));
|
||
observer
|
||
.print_changes(&[snapshot.clone()], &mut output)
|
||
.expect("observe non-prefix correction");
|
||
snapshot.response_stream = Some(response_stream(
|
||
"slot-2",
|
||
9,
|
||
1,
|
||
AGENT_RUNTIME_RESPONSE_STREAM_STATUS_READY,
|
||
"最终版",
|
||
));
|
||
observer
|
||
.print_changes(&[snapshot], &mut output)
|
||
.expect("observe new request slot");
|
||
observer
|
||
.close_response_line(&mut output)
|
||
.expect("close response line");
|
||
|
||
let output = String::from_utf8(output).expect("stream output is utf-8");
|
||
assert!(output.contains("reason=non-prefix-correction"));
|
||
assert!(output.contains("reason=new-request-slot"));
|
||
assert_eq!(output.matches("旧稿").count(), 1);
|
||
assert_eq!(output.matches("修正版").count(), 1);
|
||
assert_eq!(output.matches("最终版").count(), 1);
|
||
}
|
||
|
||
#[test]
|
||
fn response_stream_resets_sequence_for_same_run_steer_cursor() {
|
||
let mut observer = SwarmRuntimeObserver::default();
|
||
let mut output = Vec::new();
|
||
let mut initial = response_stream(
|
||
"slot-1",
|
||
9,
|
||
4,
|
||
AGENT_RUNTIME_RESPONSE_STREAM_STATUS_STREAMING,
|
||
"纠偏前回复",
|
||
);
|
||
initial.applied_steer_cursor = 1;
|
||
observer
|
||
.print_changes(&[runtime_with_response_stream(initial)], &mut output)
|
||
.expect("observe pre-steer stream");
|
||
|
||
let mut steered = response_stream(
|
||
"slot-1",
|
||
9,
|
||
1,
|
||
AGENT_RUNTIME_RESPONSE_STREAM_STATUS_READY,
|
||
"纠偏后回复",
|
||
);
|
||
steered.applied_steer_cursor = 2;
|
||
observer
|
||
.print_changes(&[runtime_with_response_stream(steered)], &mut output)
|
||
.expect("observe same-run stream after steer");
|
||
observer
|
||
.close_response_line(&mut output)
|
||
.expect("close steered response line");
|
||
|
||
let output = String::from_utf8(output).expect("steer output is utf-8");
|
||
assert!(output.contains("reason=new-steer-cursor"));
|
||
assert!(!output.contains("reason=sequence-rollback"));
|
||
assert_eq!(output.matches("纠偏前回复").count(), 1);
|
||
assert_eq!(output.matches("纠偏后回复").count(), 1);
|
||
}
|
||
|
||
#[test]
|
||
fn response_stream_reconnects_without_repeating_body_and_rejects_sequence_rollback() {
|
||
let mut observer = SwarmRuntimeObserver::default();
|
||
let mut output = Vec::new();
|
||
let mut snapshot = runtime_with_response_stream(response_stream(
|
||
"slot-1",
|
||
11,
|
||
3,
|
||
AGENT_RUNTIME_RESPONSE_STREAM_STATUS_STREAMING,
|
||
"已输出",
|
||
));
|
||
observer
|
||
.print_changes(&[snapshot.clone()], &mut output)
|
||
.expect("observe initial stream");
|
||
|
||
snapshot.response_stream = None;
|
||
assert!(observer
|
||
.print_changes(&[snapshot.clone()], &mut output)
|
||
.expect("observe disconnect"));
|
||
snapshot.response_stream = Some(response_stream(
|
||
"slot-1",
|
||
11,
|
||
3,
|
||
AGENT_RUNTIME_RESPONSE_STREAM_STATUS_STREAMING,
|
||
"已输出",
|
||
));
|
||
assert!(observer
|
||
.print_changes(&[snapshot.clone()], &mut output)
|
||
.expect("observe reconnect"));
|
||
|
||
snapshot.response_stream = Some(response_stream(
|
||
"slot-1",
|
||
11,
|
||
2,
|
||
AGENT_RUNTIME_RESPONSE_STREAM_STATUS_STREAMING,
|
||
"回退正文",
|
||
));
|
||
assert!(observer
|
||
.print_changes(&[snapshot.clone()], &mut output)
|
||
.expect("report sequence rollback"));
|
||
assert!(!observer
|
||
.print_changes(&[snapshot], &mut output)
|
||
.expect("deduplicate repeated rollback"));
|
||
|
||
let cursor = observer
|
||
.response_streams
|
||
.get("code-prototype")
|
||
.expect("response cursor");
|
||
assert_eq!(cursor.sequence, 3);
|
||
assert_eq!(cursor.accumulated_text, "已输出");
|
||
assert_eq!(cursor.printed_accumulated_text.as_deref(), Some("已输出"));
|
||
|
||
let recovered = runtime_with_response_stream(response_stream(
|
||
"slot-1",
|
||
11,
|
||
4,
|
||
AGENT_RUNTIME_RESPONSE_STREAM_STATUS_READY,
|
||
"已输出继续",
|
||
));
|
||
assert!(observer
|
||
.print_changes(&[recovered], &mut output)
|
||
.expect("resume from accepted high-water mark"));
|
||
observer
|
||
.close_response_line(&mut output)
|
||
.expect("close recovered response line");
|
||
|
||
let output = String::from_utf8(output).expect("stream output is utf-8");
|
||
assert!(output.contains("reason=reconnect"));
|
||
assert!(output.contains("reason=sequence-rollback"));
|
||
assert_eq!(output.matches("已输出").count(), 1);
|
||
assert_eq!(output.matches("继续").count(), 1);
|
||
assert!(!output.contains("回退正文"));
|
||
}
|
||
|
||
#[test]
|
||
fn settled_parent_reply_is_not_repeated_after_complete_stream() {
|
||
let mut observer = SwarmRuntimeObserver::default();
|
||
let mut output = Vec::new();
|
||
let mut snapshot = runtime_with_response_stream(response_stream(
|
||
"slot-1",
|
||
13,
|
||
4,
|
||
AGENT_RUNTIME_RESPONSE_STREAM_STATUS_READY,
|
||
"权威最终回复",
|
||
));
|
||
observer
|
||
.print_changes(&[snapshot.clone()], &mut output)
|
||
.expect("observe complete stream");
|
||
snapshot.response_stream = Some(response_stream(
|
||
"slot-1",
|
||
13,
|
||
5,
|
||
AGENT_RUNTIME_RESPONSE_STREAM_STATUS_COMMITTED,
|
||
"权威最终回复",
|
||
));
|
||
observer
|
||
.print_changes(&[snapshot], &mut output)
|
||
.expect("observe committed stream without printing body");
|
||
observer
|
||
.close_response_line(&mut output)
|
||
.expect("close response line");
|
||
print_settled_parent_reply(
|
||
"code-prototype",
|
||
"session-test",
|
||
Some("权威最终回复"),
|
||
&observer,
|
||
&mut output,
|
||
)
|
||
.expect("settle streamed reply");
|
||
let (conversation_metrics, _) =
|
||
summarize_new_assistant_messages([("assistant", "权威最终回复")]);
|
||
let report = build_swarm_turn_report(
|
||
SwarmTurnReportOutcome::Settled,
|
||
"code-prototype",
|
||
"session-test",
|
||
None,
|
||
&[],
|
||
conversation_metrics,
|
||
0,
|
||
)
|
||
.expect("build streamed reply turn report");
|
||
print_turn_outcome(SwarmTurnOutcome::Settled(report), &mut output)
|
||
.expect("print settled report after stream");
|
||
|
||
let output = String::from_utf8(output).expect("settle output is utf-8");
|
||
assert_eq!(output.matches("权威最终回复").count(), 1);
|
||
assert!(output.contains("父 Agent 回复已完整流式输出"));
|
||
assert!(output.contains(SWARM_TURN_REPORT_PREFIX));
|
||
|
||
let mut next_run_stream = response_stream(
|
||
"slot-next",
|
||
14,
|
||
1,
|
||
AGENT_RUNTIME_RESPONSE_STREAM_STATUS_READY,
|
||
"下一轮回复",
|
||
);
|
||
next_run_stream.run_id = "run-next".to_string();
|
||
let mut next_run_observer = SwarmRuntimeObserver::default();
|
||
let mut next_run_output = Vec::new();
|
||
next_run_observer
|
||
.print_changes(
|
||
&[runtime_with_response_stream(next_run_stream)],
|
||
&mut next_run_output,
|
||
)
|
||
.expect("observe next run stream");
|
||
print_settled_parent_reply_for_run(
|
||
"code-prototype",
|
||
"session-test",
|
||
Some("run-target"),
|
||
Some("上一轮最终回复"),
|
||
&next_run_observer,
|
||
&mut next_run_output,
|
||
)
|
||
.expect("next run stream must not suppress target reply");
|
||
let next_run_output = String::from_utf8(next_run_output).expect("next run output is utf-8");
|
||
assert!(next_run_output.contains("Agent> 上一轮最终回复"));
|
||
|
||
let mut fallback = Vec::new();
|
||
print_settled_parent_reply(
|
||
"code-prototype",
|
||
"session-test",
|
||
Some("未流过的权威回复"),
|
||
&SwarmRuntimeObserver::default(),
|
||
&mut fallback,
|
||
)
|
||
.expect("print authoritative fallback");
|
||
let fallback = String::from_utf8(fallback).expect("fallback output is utf-8");
|
||
assert!(fallback.contains("Agent> 未流过的权威回复"));
|
||
}
|
||
|
||
#[test]
|
||
fn response_stream_status_reports_only_status_sequence_and_char_count() {
|
||
let body = "PRIVATE_RESPONSE_BODY";
|
||
let stream = response_stream(
|
||
"slot-private",
|
||
17,
|
||
8,
|
||
AGENT_RUNTIME_RESPONSE_STREAM_STATUS_READY,
|
||
body,
|
||
);
|
||
let mut output = Vec::new();
|
||
print_runtime_response_stream_status(Some(&stream), &mut output)
|
||
.expect("print response stream status");
|
||
let output = String::from_utf8(output).expect("status output is utf-8");
|
||
|
||
assert_eq!(
|
||
output.trim(),
|
||
format!(
|
||
"[回复流] status=ready sequence=8 chars={}",
|
||
body.chars().count()
|
||
)
|
||
);
|
||
assert!(!output.contains(body));
|
||
assert!(!output.contains("slot-private"));
|
||
}
|
||
|
||
#[test]
|
||
fn input_channel_preserves_lines_and_eof() {
|
||
let (tx, rx) = mpsc::channel();
|
||
tx.send(SwarmInputEvent::Line("hello swarm".to_string()))
|
||
.expect("send line");
|
||
tx.send(SwarmInputEvent::Eof).expect("send eof");
|
||
assert_eq!(
|
||
receive_swarm_chat_line(&rx).expect("read line").as_deref(),
|
||
Some("hello swarm")
|
||
);
|
||
assert_eq!(receive_swarm_chat_line(&rx).expect("read eof"), None);
|
||
}
|
||
|
||
#[test]
|
||
fn confirmation_prompt_defers_to_bare_goal_status_without_deciding_action() {
|
||
let root = std::env::temp_dir().join(format!(
|
||
"swarm-goal-confirmation-{}-{}",
|
||
std::process::id(),
|
||
unix_millis()
|
||
));
|
||
init_local_game_project_at(&root, "project-1", "Goal 确认提示测试")
|
||
.expect("initialize Goal prompt project");
|
||
let (tx, rx) = mpsc::channel();
|
||
tx.send(SwarmInputEvent::Line("/goal".to_string()))
|
||
.expect("send bare Goal status");
|
||
let mut output = Vec::new();
|
||
|
||
let decision = prompt_swarm_decision(
|
||
&root,
|
||
GAME_CREATOR_PROJECT_SUPERVISOR_AGENT_ID,
|
||
&rx,
|
||
&mut output,
|
||
"",
|
||
)
|
||
.expect("handle Goal status during confirmation");
|
||
assert!(matches!(decision, SwarmPromptDecision::Deferred));
|
||
let output = String::from_utf8(output).expect("prompt output is utf-8");
|
||
assert!(output.contains("当前尚未设置持久目标"));
|
||
assert!(!output.contains("[已批准]"));
|
||
assert!(!output.contains("[已拒绝]"));
|
||
|
||
fs::remove_dir_all(root).ok();
|
||
}
|
||
|
||
#[test]
|
||
fn event_deduplication_keeps_phase_and_detail_changes() {
|
||
let base = serde_json::json!({
|
||
"agentId": "code-prototype",
|
||
"taskId": "task-1",
|
||
"sessionId": "session-1",
|
||
"runId": "run-1",
|
||
"eventType": "observation",
|
||
"status": "running",
|
||
"phase": "action",
|
||
"summary": "工具观察",
|
||
"detail": "第一条",
|
||
"updatedAt": 100,
|
||
});
|
||
let first =
|
||
serde_json::from_value::<AgentRuntimeEvent>(base.clone()).expect("deserialize first event");
|
||
let mut changed = base;
|
||
changed["phase"] = serde_json::json!("observation");
|
||
changed["detail"] = serde_json::json!("第二条");
|
||
let second =
|
||
serde_json::from_value::<AgentRuntimeEvent>(changed).expect("deserialize second event");
|
||
assert_ne!(runtime_event_key(&first), runtime_event_key(&second));
|
||
}
|