修复Swarm连续任务恢复与终态归属
Project CI / Repository checks (push) Successful in 2m43s
Project CI / Backend tests (push) Successful in 7m41s
Project CI / Native shell tests (push) Failing after 12m2s
Project CI / Frontend tests (push) Failing after 1m10s

以启动 mutation 返回的 acceptedRunId 建立不可变 turn baseline

拆分队列忙碌与 Runtime 可追加状态并恢复 cancelled canonical 后的 pending 任务

按 finalization message ID 和完整 task journal 隔离连续 B/C 对话、失败扫描与报告计数

以 agentId 和 runId 聚合专业 Agent 终态并保持 journal 失败优先

补齐重复投递、恢复幂等、连续轮次和历史 repair 回归及项目文档
This commit is contained in:
AIGameCreator App
2026-07-25 14:07:54 +08:00
parent 3237ad140a
commit 006846e667
17 changed files with 1764 additions and 196 deletions
@@ -278,54 +278,63 @@ pub(in crate::agent) fn resume_game_creator_agent_background_tasks_unredacted_at
}
AgentRuntimeFinalizationResume::NotFound(runtime_lock) => runtime_lock,
};
if let Some(result) =
let cancelled_with_queued_work = if let Some(result) =
reconcile_game_creator_agent_control_before_resume_at(root, &agent_id)?
{
if result.state.phase != "cancelled" {
resumed.push(result);
}
drop(runtime_lock);
continue;
}
if let Some(result) =
reconcile_game_creator_agent_process_sessions_after_restart_at(root, &agent_id)?
{
resumed.push(result);
drop(runtime_lock);
continue;
}
let runtime_lock = match resume_game_creator_agent_parallel_read_batch_at(
root,
&agent_id,
runtime_lock,
)? {
AgentRuntimePendingActionResume::Handled(result) => {
resumed.push(result);
let cancelled_with_queued_work = result.state.phase == "cancelled"
&& result.task_queue.pending > 0
&& result.task_queue.running == 0;
if !cancelled_with_queued_work {
if result.state.phase != "cancelled" {
resumed.push(result);
}
drop(runtime_lock);
continue;
}
AgentRuntimePendingActionResume::NotFound(runtime_lock) => runtime_lock,
true
} else {
false
};
let runtime_lock = match resume_game_creator_agent_pending_tool_action_at(
root,
&agent_id,
runtime_lock,
)? {
AgentRuntimePendingActionResume::Handled(result) => {
let runtime_lock = if cancelled_with_queued_work {
runtime_lock
} else {
if let Some(result) =
reconcile_game_creator_agent_process_sessions_after_restart_at(root, &agent_id)?
{
resumed.push(result);
drop(runtime_lock);
continue;
}
AgentRuntimePendingActionResume::NotFound(runtime_lock) => runtime_lock,
};
let runtime_lock = match resume_game_creator_agent_provider_action_batch_at(
root,
&agent_id,
runtime_lock,
)? {
AgentRuntimePendingActionResume::Handled(result) => {
resumed.push(result);
continue;
let runtime_lock = match resume_game_creator_agent_parallel_read_batch_at(
root,
&agent_id,
runtime_lock,
)? {
AgentRuntimePendingActionResume::Handled(result) => {
resumed.push(result);
continue;
}
AgentRuntimePendingActionResume::NotFound(runtime_lock) => runtime_lock,
};
let runtime_lock = match resume_game_creator_agent_pending_tool_action_at(
root,
&agent_id,
runtime_lock,
)? {
AgentRuntimePendingActionResume::Handled(result) => {
resumed.push(result);
continue;
}
AgentRuntimePendingActionResume::NotFound(runtime_lock) => runtime_lock,
};
match resume_game_creator_agent_provider_action_batch_at(root, &agent_id, runtime_lock)?
{
AgentRuntimePendingActionResume::Handled(result) => {
resumed.push(result);
continue;
}
AgentRuntimePendingActionResume::NotFound(runtime_lock) => runtime_lock,
}
AgentRuntimePendingActionResume::NotFound(runtime_lock) => runtime_lock,
};
let Some(task) =
read_recoverable_runnable_game_creator_agent_runtime_task(root, &agent_id)?
@@ -46,7 +46,10 @@ pub(crate) fn start_game_creator_agent_background_task_for_session_at(
start_game_creator_agent_background_task_with_run_id_for_session_at(
root, agent_id, session_id, task, run_id,
)
.map(|(result, _run_id)| result)
.map(|(mut result, accepted_run_id)| {
result.accepted_run_id = Some(accepted_run_id);
result
})
}
#[cfg(test)]
@@ -132,7 +135,10 @@ pub(crate) fn start_game_creator_supervisor_background_task_for_session_at(
source,
Some(run_profile),
)
.map(|(result, _run_id)| result)
.map(|(mut result, accepted_run_id)| {
result.accepted_run_id = Some(accepted_run_id);
result
})
}
pub(in crate::agent) fn start_game_creator_agent_background_task_with_source_at(
@@ -42,6 +42,7 @@ pub(crate) use context_window::{
};
#[allow(unused_imports)]
pub(crate) use finalization::{
game_creator_agent_runtime_finalization_message_id,
game_creator_agent_runtime_finalization_path,
read_game_creator_agent_runtime_finalization_journal,
remove_game_creator_agent_runtime_finalization_journal,
@@ -77,7 +78,8 @@ pub(crate) use run_configuration::{
read_game_creator_agent_runtime_run_profile_binding,
};
pub(crate) use steering::{
consume_game_creator_agent_runtime_steers, game_creator_agent_runtime_steer_ledger_path,
consume_game_creator_agent_runtime_steers, game_creator_agent_runtime_accepts_steer,
game_creator_agent_runtime_steer_ledger_path,
interrupt_game_creator_agent_runtime_provider_request_at,
render_game_creator_agent_runtime_steers_for_prompt, steer_game_creator_agent_runtime_task_at,
steer_game_creator_agent_runtime_task_for_profile_at,
@@ -21,7 +21,7 @@ pub(crate) fn game_creator_agent_runtime_finalization_path(
))
}
pub(in crate::agent) fn game_creator_agent_runtime_finalization_message_id(
pub(crate) fn game_creator_agent_runtime_finalization_message_id(
agent_id: &str,
session_id: &str,
run_id: &str,
@@ -429,9 +429,36 @@ pub(in crate::agent) fn ensure_game_creator_agent_runtime_steer_audit(
)
}
pub(crate) fn game_creator_agent_runtime_accepts_steer(state: &AgentRuntimeState) -> bool {
!state.agent_id.starts_with("child-")
&& state.source != AGENT_RUNTIME_ISOLATED_CHILD_SOURCE
&& state.status != "waiting-for-user-input"
&& state.phase != "waiting-for-user-input"
&& !matches!(
state.status.as_str(),
"completed" | "failed" | "cancelled" | "cancelling"
)
&& !matches!(
state.phase.as_str(),
"completed"
| "failed"
| "cancelled"
| "cancelling"
| "finalizing"
| "needs-reconciliation"
)
&& matches!(
state.status.as_str(),
"running" | "waiting-for-confirmation"
)
}
pub(in crate::agent) fn validate_agent_runtime_steer_target_state(
state: &AgentRuntimeState,
) -> Result<(), String> {
if game_creator_agent_runtime_accepts_steer(state) {
return Ok(());
}
if state.agent_id.starts_with("child-") || state.source == AGENT_RUNTIME_ISOLATED_CHILD_SOURCE {
return Err("动态隔离子 Agent 暂不接受运行中追加指令".to_string());
}
@@ -462,7 +489,7 @@ pub(in crate::agent) fn validate_agent_runtime_steer_target_state(
state.status
));
}
Ok(())
unreachable!("可追加状态已在函数入口返回")
}
pub(crate) fn steer_game_creator_agent_runtime_task_at(
@@ -44,7 +44,7 @@ pub(crate) use delegation::{
};
pub(crate) use delivery::{
agent_runtime_delegation_id, dispatch_isolated_agent_join_at,
publish_game_creator_agent_delegate_result,
game_creator_agent_runtime_terminal_status, publish_game_creator_agent_delegate_result,
};
#[allow(unused_imports)]
pub(crate) use isolated_joins::{
@@ -37,7 +37,7 @@ pub(in crate::agent) fn agent_runtime_delegate_receipt_run_id(delegation_id: &st
)
}
pub(in crate::agent) fn game_creator_agent_runtime_terminal_status(
pub(crate) fn game_creator_agent_runtime_terminal_status(
task: &AgentRuntimeTaskRecord,
) -> Option<&'static str> {
match task.phase.as_str() {
@@ -153,12 +153,37 @@ pub(super) fn read_turn_conversation_snapshot(
if baseline.previous_message_count > conversation.messages.len() {
return Err("Swarm turn 对话 baseline 超出当前 Session 消息数".to_string());
}
let (mut metrics, final_reply) = summarize_new_assistant_messages(
let target_message_id = game_creator_agent_runtime_finalization_message_id(
parent_agent_id,
session_id,
&baseline.parent_run_id,
);
if let Some(journal) = read_game_creator_agent_runtime_finalization_journal(
root,
parent_agent_id,
&baseline.parent_run_id,
)? {
if journal.agent_id != parent_agent_id
|| journal.session_id != session_id
|| journal.run_id != baseline.parent_run_id
|| journal.message_id != target_message_id
{
return Err("Swarm finalization 与目标 parent Session/run 不匹配".to_string());
}
}
let (mut metrics, final_reply) = summarize_scoped_new_assistant_messages(
conversation
.messages
.iter()
.skip(baseline.previous_message_count)
.map(|message| (message.role.as_str(), message.content.as_str())),
.map(|message| {
(
message.role.as_str(),
message.content.as_str(),
message.message_id.as_deref(),
)
}),
Some(&target_message_id),
);
let mut final_reply = final_reply.map(str::to_string);
let recovered_before_observation = baseline.recovered_assistant.is_some();
@@ -198,14 +223,46 @@ pub(super) fn read_turn_conversation_metrics(
pub(super) fn summarize_new_assistant_messages<'a>(
messages: impl IntoIterator<Item = (&'a str, &'a str)>,
) -> (SwarmTurnConversationMetrics, Option<&'a str>) {
let mut new_assistant_message_count = 0;
let mut final_reply = None;
for (role, content) in messages {
if role == "assistant" {
new_assistant_message_count += 1;
final_reply = Some(content);
summarize_scoped_new_assistant_messages(
messages
.into_iter()
.map(|(role, content)| (role, content, None)),
None,
)
}
pub(super) fn summarize_scoped_new_assistant_messages<'a>(
messages: impl IntoIterator<Item = (&'a str, &'a str, Option<&'a str>)>,
target_message_id: Option<&str>,
) -> (SwarmTurnConversationMetrics, Option<&'a str>) {
let mut scoped_count = 0;
let mut scoped_reply = None;
let mut legacy_count = 0;
let mut legacy_reply = None;
let mut identified_assistant_exists = false;
for (role, content, message_id) in messages {
if role != "assistant" {
continue;
}
match message_id {
Some(message_id) if target_message_id == Some(message_id) => {
identified_assistant_exists = true;
scoped_count += 1;
scoped_reply = Some(content);
}
None => {
legacy_count += 1;
legacy_reply = Some(content);
}
Some(_) => identified_assistant_exists = true,
}
}
let (new_assistant_message_count, final_reply) =
if target_message_id.is_some() && (scoped_count > 0 || identified_assistant_exists) {
(scoped_count, scoped_reply)
} else {
(legacy_count, legacy_reply)
};
(
SwarmTurnConversationMetrics {
new_assistant_message_count,
@@ -229,9 +286,10 @@ pub(super) fn print_new_parent_reply<W: Write>(
writeln!(output, "[本轮结束] 父 Agent 回复已在恢复前持久化。")
.map_err(|error| format!("写入终端失败:{error}"))?;
} else {
print_settled_parent_reply(
print_settled_parent_reply_for_run(
parent_agent_id,
session_id,
Some(&baseline.parent_run_id),
snapshot.final_reply.as_deref(),
observer,
output,
@@ -246,12 +304,31 @@ pub(super) fn print_settled_parent_reply<W: Write>(
reply: Option<&str>,
observer: &SwarmRuntimeObserver,
output: &mut W,
) -> Result<(), String> {
print_settled_parent_reply_for_run(parent_agent_id, session_id, None, reply, observer, output)
}
pub(super) fn print_settled_parent_reply_for_run<W: Write>(
parent_agent_id: &str,
session_id: &str,
target_run_id: Option<&str>,
reply: Option<&str>,
observer: &SwarmRuntimeObserver,
output: &mut W,
) -> Result<(), String> {
let Some(reply) = reply else {
return writeln!(output, "[本轮结束] 父 Agent 未产生新的最终回复。")
.map_err(|error| format!("写入终端失败:{error}"));
};
if observer.parent_reply_was_fully_streamed(parent_agent_id, session_id, reply) {
let stream_belongs_to_target = target_run_id.is_none_or(|target_run_id| {
observer
.response_streams
.get(parent_agent_id)
.is_some_and(|cursor| cursor.identity.run_id == target_run_id)
});
if stream_belongs_to_target
&& observer.parent_reply_was_fully_streamed(parent_agent_id, session_id, reply)
{
writeln!(output, "[本轮结束] 父 Agent 回复已完整流式输出。")
.map_err(|error| format!("写入终端失败:{error}"))
} else {
@@ -52,53 +52,189 @@ pub(super) struct SwarmTurnReport {
pub(super) reconciliation_agent_count: usize,
}
#[derive(Default)]
struct SwarmTurnRuntimeMetrics {
runtime_count: usize,
busy_runtime_count: usize,
pending_task_count: u64,
running_task_count: u64,
waiting_for_confirmation_count: u64,
waiting_for_user_input_count: u64,
}
struct SwarmTurnRuntimeSnapshot {
agent_id: String,
run_id: String,
status: String,
phase: String,
updated_at: u64,
}
fn runtime_state_belongs_to_turn(
state: &AgentRuntimeState,
parent_agent_id: &str,
session_id: &str,
parent_run_id: &str,
) -> bool {
(state.agent_id == parent_agent_id
&& state.session_id == session_id
&& state.run_id == parent_run_id)
|| (state.parent_agent_id.as_deref() == Some(parent_agent_id)
&& state.parent_run_id.as_deref() == Some(parent_run_id))
}
fn runtime_task_belongs_to_turn(
task: &AgentRuntimeTaskRecord,
parent_agent_id: &str,
session_id: &str,
parent_run_id: &str,
) -> bool {
(task.agent_id == parent_agent_id
&& task.session_id == session_id
&& task.run_id == parent_run_id)
|| (task.parent_agent_id.as_deref() == Some(parent_agent_id)
&& task.parent_run_id.as_deref() == Some(parent_run_id))
}
fn upsert_turn_runtime_snapshot(
snapshots: &mut Vec<SwarmTurnRuntimeSnapshot>,
candidate: SwarmTurnRuntimeSnapshot,
) {
if candidate.run_id.trim().is_empty() {
return;
}
if let Some(existing) = snapshots.iter_mut().find(|snapshot| {
snapshot.agent_id == candidate.agent_id && snapshot.run_id == candidate.run_id
}) {
if candidate.updated_at >= existing.updated_at {
*existing = candidate;
}
} else {
snapshots.push(candidate);
}
}
fn scoped_swarm_turn_runtime_metrics(
parent_agent_id: &str,
session_id: &str,
parent_run_id: Option<&str>,
runtimes: &[AgentRuntimeResult],
) -> Result<SwarmTurnRuntimeMetrics, String> {
let Some(parent_run_id) = parent_run_id else {
return Ok(SwarmTurnRuntimeMetrics::default());
};
let mut snapshots = Vec::new();
for runtime in runtimes {
let journal_tasks = if runtime.task_path.trim().is_empty() {
None
} else {
Some(read_all_game_creator_agent_runtime_tasks(Path::new(
&runtime.task_path,
))?)
};
let tasks = journal_tasks.as_deref().unwrap_or(&runtime.recent_tasks);
for task in tasks {
if runtime_task_belongs_to_turn(task, parent_agent_id, session_id, parent_run_id) {
upsert_turn_runtime_snapshot(
&mut snapshots,
SwarmTurnRuntimeSnapshot {
agent_id: task.agent_id.clone(),
run_id: task.run_id.clone(),
status: task.status.clone(),
phase: task.phase.clone(),
updated_at: task.updated_at,
},
);
}
}
if runtime_state_belongs_to_turn(&runtime.state, parent_agent_id, session_id, parent_run_id)
{
upsert_turn_runtime_snapshot(
&mut snapshots,
SwarmTurnRuntimeSnapshot {
agent_id: runtime.state.agent_id.clone(),
run_id: runtime.state.run_id.clone(),
status: runtime.state.status.clone(),
phase: runtime.state.phase.clone(),
updated_at: runtime.state.updated_at,
},
);
}
}
let mut metrics = SwarmTurnRuntimeMetrics {
runtime_count: snapshots.len(),
..SwarmTurnRuntimeMetrics::default()
};
for snapshot in snapshots {
if matches!(
snapshot.status.as_str(),
"pending"
| "running"
| "waiting-for-confirmation"
| "waiting-for-user-input"
| "cancelling"
) || snapshot.phase == "needs-reconciliation"
{
metrics.busy_runtime_count += 1;
}
match snapshot.status.as_str() {
"pending" => metrics.pending_task_count += 1,
"running" => metrics.running_task_count += 1,
"waiting-for-confirmation" => metrics.waiting_for_confirmation_count += 1,
"waiting-for-user-input" => metrics.waiting_for_user_input_count += 1,
_ => {}
}
}
Ok(metrics)
}
pub(super) fn build_swarm_turn_report(
outcome: SwarmTurnReportOutcome,
parent_agent_id: &str,
session_id: &str,
expected_parent_run_id: Option<&str>,
runtimes: &[AgentRuntimeResult],
conversation_metrics: SwarmTurnConversationMetrics,
reconciliation_agent_count: usize,
) -> SwarmTurnReport {
let parent_run_id = runtimes
.iter()
.find(|runtime| {
runtime.state.agent_id == parent_agent_id && runtime.state.session_id == session_id
})
.map(|runtime| runtime.state.run_id.trim())
) -> Result<SwarmTurnReport, String> {
let parent_run_id = expected_parent_run_id
.map(str::trim)
.filter(|run_id| !run_id.is_empty())
.map(str::to_string);
SwarmTurnReport {
.map(str::to_string)
.or_else(|| {
runtimes
.iter()
.find(|runtime| {
runtime.state.agent_id == parent_agent_id
&& runtime.state.session_id == session_id
})
.map(|runtime| runtime.state.run_id.trim())
.filter(|run_id| !run_id.is_empty())
.map(str::to_string)
});
let runtime_metrics = scoped_swarm_turn_runtime_metrics(
parent_agent_id,
session_id,
parent_run_id.as_deref(),
runtimes,
)?;
Ok(SwarmTurnReport {
schema_version: SWARM_TURN_REPORT_SCHEMA_VERSION,
outcome,
parent_agent_id: parent_agent_id.to_string(),
session_id: session_id.to_string(),
parent_run_id,
runtime_count: runtimes.len(),
busy_runtime_count: runtimes
.iter()
.filter(|runtime| runtime_is_busy(runtime))
.count(),
pending_task_count: runtimes
.iter()
.map(|runtime| u64::from(runtime.task_queue.pending))
.sum(),
running_task_count: runtimes
.iter()
.map(|runtime| u64::from(runtime.task_queue.running))
.sum(),
waiting_for_confirmation_count: runtimes
.iter()
.map(|runtime| u64::from(runtime.task_queue.waiting_for_confirmation))
.sum(),
waiting_for_user_input_count: runtimes
.iter()
.map(|runtime| u64::from(runtime.task_queue.waiting_for_user_input))
.sum(),
runtime_count: runtime_metrics.runtime_count,
busy_runtime_count: runtime_metrics.busy_runtime_count,
pending_task_count: runtime_metrics.pending_task_count,
running_task_count: runtime_metrics.running_task_count,
waiting_for_confirmation_count: runtime_metrics.waiting_for_confirmation_count,
waiting_for_user_input_count: runtime_metrics.waiting_for_user_input_count,
new_assistant_message_count: conversation_metrics.new_assistant_message_count,
final_reply_chars: conversation_metrics.final_reply_chars,
reconciliation_agent_count,
}
})
}
pub(super) fn print_turn_outcome<W: Write>(
@@ -34,6 +34,205 @@ pub(super) fn swarm_parent_runtime<'a>(
})
}
pub(super) fn swarm_parent_runtime_for_run<'a>(
parent_agent_id: &str,
session_id: &str,
run_profile: &str,
expected_run_id: &str,
runtimes: &'a [AgentRuntimeResult],
) -> Option<&'a AgentRuntimeResult> {
swarm_parent_runtime(parent_agent_id, session_id, run_profile, runtimes)
.filter(|runtime| runtime.state.run_id == expected_run_id)
}
pub(super) fn swarm_parent_steer_target<'a>(
parent_agent_id: &str,
session_id: &str,
run_profile: &str,
expected_run_id: Option<&str>,
runtimes: &'a [AgentRuntimeResult],
) -> Option<&'a AgentRuntimeResult> {
swarm_parent_runtime(parent_agent_id, session_id, run_profile, runtimes).filter(|runtime| {
expected_run_id.is_none_or(|run_id| runtime.state.run_id == run_id)
&& game_creator_agent_runtime_accepts_steer(&runtime.state)
})
}
pub(super) fn matching_pending_swarm_run_id<'a>(
runtime: &'a AgentRuntimeResult,
session_id: &str,
run_profile: &str,
message: &str,
) -> Option<&'a str> {
let message = message.trim();
runtime
.recent_tasks
.iter()
.rev()
.find(|task| {
task.session_id == session_id
&& task.run_profile == run_profile
&& task.status == "pending"
&& task.phase == "queued"
&& task.task.trim() == message
})
.map(|task| task.run_id.as_str())
}
pub(super) fn next_pending_swarm_task<'a>(
runtime: &'a AgentRuntimeResult,
session_id: &str,
run_profile: &str,
) -> Option<&'a AgentRuntimeTaskRecord> {
runtime.recent_tasks.iter().find(|task| {
task.session_id == session_id
&& task.run_profile == run_profile
&& task.status == "pending"
&& task.phase == "queued"
})
}
pub(super) fn swarm_parent_task_for_run<'a>(
parent_agent_id: &str,
session_id: &str,
run_profile: &str,
expected_run_id: &str,
runtimes: &'a [AgentRuntimeResult],
) -> Option<&'a AgentRuntimeTaskRecord> {
swarm_parent_runtime(parent_agent_id, session_id, run_profile, runtimes).and_then(|runtime| {
runtime.recent_tasks.iter().find(|task| {
task.agent_id == parent_agent_id
&& task.session_id == session_id
&& task.run_profile == run_profile
&& task.run_id == expected_run_id
})
})
}
pub(super) fn runtime_result_from_task_record(task: &AgentRuntimeTaskRecord) -> AgentRuntimeResult {
let mut state = default_game_creator_agent_runtime_state(&task.agent_id, &task.run_id);
state.task_id = task.task_id.clone();
state.session_id = task.session_id.clone();
state.source = task.source.clone();
state.run_profile = task.run_profile.clone();
state.run_profile_binding_fingerprint = task.run_profile_binding_fingerprint.clone();
state.parent_agent_id = task.parent_agent_id.clone();
state.parent_run_id = task.parent_run_id.clone();
state.delegation_id = task.delegation_id.clone();
state.goal_id = task.goal_id.clone();
state.goal_revision = task.goal_revision;
state.goal_status = task.goal_status.clone();
state.current_task = task.task.clone();
state.status = task.status.clone();
state.phase = task.phase.clone();
state.current_action = task.current_action.clone();
state.error = task.error.clone();
state.updated_at = task.updated_at;
AgentRuntimeResult {
state,
accepted_run_id: None,
session_path: String::new(),
event_path: String::new(),
task_path: String::new(),
task_queue: AgentRuntimeTaskQueueSummary::default(),
recent_events: Vec::new(),
recent_tasks: vec![task.clone()],
response_stream: None,
user_input_request: None,
}
}
pub(super) fn swarm_parent_runtime_snapshot_for_run(
root: &Path,
parent_agent_id: &str,
session_id: &str,
run_profile: &str,
expected_run_id: &str,
runtimes: &[AgentRuntimeResult],
) -> Result<Option<AgentRuntimeResult>, String> {
if let Some(runtime) = swarm_parent_runtime_for_run(
parent_agent_id,
session_id,
run_profile,
expected_run_id,
runtimes,
) {
return Ok(Some(runtime.clone()));
}
if let Some(task) = read_latest_game_creator_agent_runtime_task_by_run_id(
root,
parent_agent_id,
expected_run_id,
)? {
if task.agent_id != parent_agent_id
|| task.session_id != session_id
|| task.run_profile != run_profile
|| task.run_id != expected_run_id
{
return Err("Swarm parent task journal 与目标 Session/run 不匹配".to_string());
}
return Ok(Some(runtime_result_from_task_record(&task)));
}
Ok(swarm_parent_task_for_run(
parent_agent_id,
session_id,
run_profile,
expected_run_id,
runtimes,
)
.map(runtime_result_from_task_record))
}
pub(super) fn swarm_turn_is_busy(
root: &Path,
parent_agent_id: &str,
session_id: &str,
run_profile: &str,
expected_run_id: &str,
runtimes: &[AgentRuntimeResult],
) -> Result<bool, String> {
let Some(parent) = swarm_parent_runtime_snapshot_for_run(
root,
parent_agent_id,
session_id,
run_profile,
expected_run_id,
runtimes,
)?
else {
return Ok(false);
};
Ok(matches!(
parent.state.status.as_str(),
"pending"
| "running"
| "waiting-for-confirmation"
| "waiting-for-user-input"
| "cancelling"
) || parent.state.phase == "needs-reconciliation")
}
pub(super) fn swarm_current_runtimes_for_run(
parent_agent_id: &str,
session_id: &str,
run_profile: &str,
expected_run_id: &str,
runtimes: &[AgentRuntimeResult],
) -> Vec<AgentRuntimeResult> {
runtimes
.iter()
.filter(|runtime| {
(runtime.state.agent_id == parent_agent_id
&& runtime.state.session_id == session_id
&& runtime.state.run_profile == run_profile
&& runtime.state.run_id == expected_run_id)
|| (runtime.state.parent_agent_id.as_deref() == Some(parent_agent_id)
&& runtime.state.parent_run_id.as_deref() == Some(expected_run_id))
})
.cloned()
.collect()
}
pub(super) fn runtime_terminal_failure_kind(runtime: &AgentRuntimeResult) -> Option<&'static str> {
if runtime.state.phase == "needs-reconciliation" {
None
@@ -41,7 +240,9 @@ pub(super) fn runtime_terminal_failure_kind(runtime: &AgentRuntimeResult) -> Opt
Some("budget-exhausted")
} else if runtime.state.status == "cancelled" || runtime.state.phase == "cancelled" {
Some("cancelled")
} else if runtime.state.status == "failed" {
} else if runtime.state.phase == "completion-contract-failed" {
Some("completion-contract-failed")
} else if runtime.state.status == "failed" || runtime.state.phase == "failed" {
Some("failed")
} else {
None
@@ -135,13 +336,58 @@ pub(super) fn scan_swarm_terminal_failures_at(
parent_agent_id: &str,
session_id: &str,
run_profile: &str,
expected_parent_run_id: &str,
runtimes: &[AgentRuntimeResult],
) -> SwarmTerminalFailureScan {
let mut scan = SwarmTerminalFailureScan::default();
let Some(parent) = swarm_parent_runtime(parent_agent_id, session_id, run_profile, runtimes)
else {
return scan;
let canonical_parent = swarm_parent_runtime_for_run(
parent_agent_id,
session_id,
run_profile,
expected_parent_run_id,
runtimes,
);
let journal_parent_task = match read_latest_game_creator_agent_runtime_task_by_run_id(
root,
parent_agent_id,
expected_parent_run_id,
) {
Ok(Some(task))
if task.agent_id == parent_agent_id
&& task.session_id == session_id
&& task.run_profile == run_profile
&& task.run_id == expected_parent_run_id =>
{
Some(task)
}
Ok(Some(_)) => {
scan.reconciliation_agents.push(parent_agent_id.to_string());
return scan;
}
Ok(None) => None,
Err(_) => {
scan.reconciliation_agents.push(parent_agent_id.to_string());
return scan;
}
};
let parent_task = journal_parent_task.as_ref().or_else(|| {
swarm_parent_task_for_run(
parent_agent_id,
session_id,
run_profile,
expected_parent_run_id,
runtimes,
)
});
if canonical_parent.is_none()
&& parent_task.is_none_or(|task| game_creator_agent_runtime_terminal_status(task).is_none())
{
return scan;
}
let parent = canonical_parent
.cloned()
.or_else(|| parent_task.map(runtime_result_from_task_record))
.expect("target parent runtime or terminal task exists");
let claimed_deliveries = match claimed_static_delegate_deliveries_at(
root,
&parent.state.agent_id,
@@ -154,16 +400,79 @@ pub(super) fn scan_swarm_terminal_failures_at(
return scan;
}
};
if let Some(kind) = runtime_terminal_failure_kind(parent) {
if let Some(kind) = runtime_terminal_failure_kind(&parent) {
scan.failed_agents
.push(format!("{}:{kind}", parent.state.agent_id));
}
for child in runtimes.iter().filter(|runtime| {
let task_matches_parent = |task: &AgentRuntimeTaskRecord| {
task.source == "agent-delegate"
&& task.parent_agent_id.as_deref() == Some(parent.state.agent_id.as_str())
&& task.parent_run_id.as_deref() == Some(parent.state.run_id.as_str())
};
let mut specialist_agent_ids = runtimes
.iter()
.map(|runtime| runtime.state.agent_id.clone())
.filter(|agent_id| agent_id != &parent.state.agent_id)
.collect::<BTreeSet<_>>();
specialist_agent_ids.extend(
claimed_deliveries
.iter()
.map(|delivery| delivery.target_agent_id.clone()),
);
// recent_tasks remains a compatibility fallback for in-memory callers, while every
// discoverable specialist journal below overwrites it with append-order latest records.
let mut latest_child_tasks_by_identity =
BTreeMap::<(String, String), AgentRuntimeTaskRecord>::new();
for runtime in runtimes {
for task in runtime
.recent_tasks
.iter()
.filter(|task| task_matches_parent(task))
{
latest_child_tasks_by_identity
.insert((task.agent_id.clone(), task.run_id.clone()), task.clone());
}
}
for agent_id in specialist_agent_ids {
let path = game_creator_agent_runtime_task_path(root, &agent_id);
let tasks = match read_all_game_creator_agent_runtime_tasks(&path) {
Ok(tasks) => tasks,
Err(_) => {
scan.reconciliation_agents.push(agent_id);
continue;
}
};
for task in tasks.into_iter().filter(|task| task_matches_parent(task)) {
if task.agent_id != agent_id {
scan.reconciliation_agents.push(agent_id.clone());
continue;
}
latest_child_tasks_by_identity
.insert((task.agent_id.clone(), task.run_id.clone()), task);
}
}
let mut failed_children_by_identity = latest_child_tasks_by_identity
.into_iter()
.filter_map(|(identity, task)| {
let runtime = runtime_result_from_task_record(&task);
runtime_terminal_failure_kind(&runtime).map(|_| (identity, runtime))
})
.collect::<BTreeMap<_, _>>();
for runtime in runtimes.iter().filter(|runtime| {
runtime.state.source == "agent-delegate"
&& runtime.state.parent_agent_id.as_deref() == Some(parent.state.agent_id.as_str())
&& runtime.state.parent_run_id.as_deref() == Some(parent.state.run_id.as_str())
&& runtime_terminal_failure_kind(runtime).is_some()
}) {
if runtime_terminal_failure_kind(runtime).is_some() {
failed_children_by_identity.insert(
(runtime.state.agent_id.clone(), runtime.state.run_id.clone()),
runtime.clone(),
);
}
}
for child in failed_children_by_identity.values() {
let Some(delegation_id) = child
.state
.delegation_id
@@ -187,7 +496,7 @@ pub(super) fn scan_swarm_terminal_failures_at(
};
let successful_repair =
original_delivery_has_successful_repair(&delivery, &claimed_deliveries);
match classify_failed_specialist(parent, child, Some(&delivery), successful_repair) {
match classify_failed_specialist(&parent, child, Some(&delivery), successful_repair) {
SwarmSpecialistFailureDisposition::Recoverable => {}
SwarmSpecialistFailureDisposition::Failed => scan.failed_agents.push(format!(
"{}:{}",
@@ -388,10 +697,11 @@ pub(super) fn build_reconciliation_turn_outcome(
SwarmTurnReportOutcome::NeedsReconciliation,
parent_agent_id,
session_id,
Some(&conversation_baseline.parent_run_id),
runtimes,
conversation_metrics,
agent_ids.len(),
);
)?;
Ok(SwarmTurnOutcome::NeedsReconciliation { agent_ids, report })
}
@@ -411,10 +721,11 @@ pub(super) fn build_failed_turn_outcome(
SwarmTurnReportOutcome::Failed,
parent_agent_id,
session_id,
Some(&conversation_baseline.parent_run_id),
runtimes,
conversation_metrics,
0,
);
)?;
Ok(SwarmTurnOutcome::Failed { agent_ids, report })
}
@@ -437,9 +748,10 @@ pub(super) fn build_incomplete_turn_outcome(
SwarmTurnReportOutcome::Incomplete,
parent_agent_id,
session_id,
Some(&conversation_baseline.parent_run_id),
runtimes,
conversation_metrics,
0,
);
)?;
Ok(SwarmTurnOutcome::Incomplete { reasons, report })
}
File diff suppressed because it is too large Load Diff
@@ -12,52 +12,46 @@ pub(super) fn handle_swarm_user_turn<W: Write>(
) -> Result<SwarmChatFlow, String> {
let before =
read_local_conversation_for_session_at(root, Some(parent_agent_id), Some(session_id))?;
if let Some(goal) = read_game_creator_agent_goal_at(root, parent_agent_id, session_id)? {
match goal.status.as_str() {
AGENT_GOAL_STATUS_ACTIVE => {
return steer_and_wait_for_swarm_turn(
root,
parent_agent_id,
session_id,
run_profile,
&goal.run_id,
message,
before.messages.len(),
input,
output,
"Goal 已追加",
);
let active_goal_run_id =
if let Some(goal) = read_game_creator_agent_goal_at(root, parent_agent_id, session_id)? {
match goal.status.as_str() {
AGENT_GOAL_STATUS_ACTIVE => Some(goal.run_id),
AGENT_GOAL_STATUS_PAUSE_REQUESTED | AGENT_GOAL_STATUS_PAUSED => {
print_swarm_goal_error(output, "当前 Goal 已暂停;请先输入 /goal resume。")?;
return Ok(SwarmChatFlow::Continue);
}
AGENT_GOAL_STATUS_CLEARING => {
print_swarm_goal_error(output, "当前 Goal 正在清理,暂不接受新消息。")?;
return Ok(SwarmChatFlow::Continue);
}
AGENT_GOAL_STATUS_NEEDS_RECONCILIATION => {
print_swarm_goal_error(
output,
"当前 Goal 需要人工 reconciliation,暂不接受新消息。",
)?;
return Ok(SwarmChatFlow::Continue);
}
AGENT_GOAL_STATUS_COMPLETED | AGENT_GOAL_STATUS_CLEARED => None,
status => {
print_swarm_goal_error(
output,
&format!("当前 Goal 状态未知,已阻止发送:{status}"),
)?;
return Ok(SwarmChatFlow::Continue);
}
}
AGENT_GOAL_STATUS_PAUSE_REQUESTED | AGENT_GOAL_STATUS_PAUSED => {
print_swarm_goal_error(output, "当前 Goal 已暂停;请先输入 /goal resume。")?;
return Ok(SwarmChatFlow::Continue);
}
AGENT_GOAL_STATUS_CLEARING => {
print_swarm_goal_error(output, "当前 Goal 正在清理,暂不接受新消息。")?;
return Ok(SwarmChatFlow::Continue);
}
AGENT_GOAL_STATUS_NEEDS_RECONCILIATION => {
print_swarm_goal_error(
output,
"当前 Goal 需要人工 reconciliation,暂不接受新消息。",
)?;
return Ok(SwarmChatFlow::Continue);
}
AGENT_GOAL_STATUS_COMPLETED | AGENT_GOAL_STATUS_CLEARED => {}
status => {
print_swarm_goal_error(
output,
&format!("当前 Goal 状态未知,已阻止发送:{status}"),
)?;
return Ok(SwarmChatFlow::Continue);
}
}
}
} else {
None
};
let runtimes = read_game_creator_agent_runtimes_at(root)?;
if let Some(runtime) = swarm_parent_runtime(parent_agent_id, session_id, run_profile, &runtimes)
.filter(|runtime| runtime_is_busy(runtime))
{
let mut runtimes = read_game_creator_agent_runtimes_at(root)?;
if let Some(runtime) = swarm_parent_steer_target(
parent_agent_id,
session_id,
run_profile,
active_goal_run_id.as_deref(),
&runtimes,
) {
return steer_and_wait_for_swarm_turn(
root,
parent_agent_id,
@@ -68,19 +62,102 @@ pub(super) fn handle_swarm_user_turn<W: Write>(
before.messages.len(),
input,
output,
"运行中输入已排队",
if active_goal_run_id.is_some() {
"Goal 已追加"
} else {
"运行中输入已排队"
},
);
}
let pending_message_is_latest = before
.messages
.last()
.is_some_and(|item| item.role == "user" && item.content.trim() == message.trim());
let matching_pending_run_id = pending_message_is_latest
.then(|| swarm_parent_runtime(parent_agent_id, session_id, run_profile, &runtimes))
.flatten()
.and_then(|runtime| {
matching_pending_swarm_run_id(runtime, session_id, run_profile, message)
})
.map(str::to_string);
let parent_has_queued_work =
swarm_parent_runtime(parent_agent_id, session_id, run_profile, &runtimes)
.is_some_and(runtime_is_busy);
if parent_has_queued_work {
require_external_agent_runner_for_cli_runtime_write(root)?;
resume_game_creator_agent_background_tasks_at(root)?;
runtimes = read_game_creator_agent_runtimes_at(root)?;
if let Some(run_id) = matching_pending_run_id {
writeln!(
output,
"[恢复] 该消息已在 run={run_id} 落盘,继续观察原任务,不重复追加。"
)
.map_err(|error| format!("写入终端失败:{error}"))?;
let conversation_baseline = new_swarm_turn_conversation_baseline(
before.messages.len().saturating_sub(1),
&run_id,
);
return wait_and_print_swarm_turn(
root,
parent_agent_id,
session_id,
run_profile,
conversation_baseline,
input,
output,
);
}
if let Some(runtime) = swarm_parent_steer_target(
parent_agent_id,
session_id,
run_profile,
active_goal_run_id.as_deref(),
&runtimes,
) {
return steer_and_wait_for_swarm_turn(
root,
parent_agent_id,
session_id,
run_profile,
&runtime.state.run_id,
message,
before.messages.len(),
input,
output,
if active_goal_run_id.is_some() {
"Goal 已恢复并追加"
} else {
"排队任务已恢复,输入已追加"
},
);
}
}
if let Some(goal_run_id) = active_goal_run_id {
print_swarm_goal_error(
output,
&format!(
"当前 Goal run={goal_run_id} 没有可追加的运行态;请先输入 /resume 检查恢复结果。"
),
)?;
return Ok(SwarmChatFlow::Continue);
}
let action = if game_creator_agent_uses_interaction_kernel(parent_agent_id) {
let Some(runtime_lock) =
try_acquire_game_creator_agent_runtime_task_lock(root, parent_agent_id)?
else {
let current_runtimes = read_game_creator_agent_runtimes_at(root)?;
if let Some(runtime) =
swarm_parent_runtime(parent_agent_id, session_id, run_profile, &current_runtimes)
.filter(|runtime| runtime_is_busy(runtime))
{
if let Some(runtime) = swarm_parent_steer_target(
parent_agent_id,
session_id,
run_profile,
None,
&current_runtimes,
) {
return steer_and_wait_for_swarm_turn(
root,
parent_agent_id,
@@ -322,14 +399,15 @@ fn start_and_wait_for_swarm_turn<W: Write>(
requested_run_id.clone(),
)?,
};
let accepted_run_id = accepted_swarm_run_id(&started, &requested_run_id);
writeln!(
output,
"[已投递] agent={} session={} run={}",
parent_agent_id, started.state.session_id, requested_run_id
parent_agent_id, started.state.session_id, accepted_run_id
)
.map_err(|error| format!("写入终端失败:{error}"))?;
let conversation_baseline =
new_swarm_turn_conversation_baseline(previous_message_count, &started.state.run_id);
new_swarm_turn_conversation_baseline(previous_message_count, accepted_run_id);
wait_and_print_swarm_turn(
root,
parent_agent_id,
@@ -341,6 +419,19 @@ fn start_and_wait_for_swarm_turn<W: Write>(
)
}
pub(super) fn accepted_swarm_run_id(
started: &AgentRuntimeResult,
requested_run_id: &str,
) -> String {
started
.accepted_run_id
.as_deref()
.map(str::trim)
.filter(|run_id| !run_id.is_empty())
.unwrap_or(requested_run_id)
.to_string()
}
pub(super) fn handle_swarm_resume_turn<W: Write>(
root: &Path,
parent_agent_id: &str,
@@ -359,20 +450,40 @@ pub(super) fn handle_swarm_resume_turn<W: Write>(
let current_runtimes = read_game_creator_agent_runtimes_at(root)?;
let Some(parent) =
swarm_parent_runtime(parent_agent_id, session_id, run_profile, &current_runtimes)
.filter(|runtime| runtime_is_busy(runtime))
else {
writeln!(output, "[恢复] 当前 Session 没有可恢复的运行任务。")
.map_err(|error| format!("写入终端失败:{error}"))?;
return Ok(SwarmChatFlow::Continue);
};
let pending = next_pending_swarm_task(parent, session_id, run_profile);
let (target_run_id, target_status, target_phase) =
if game_creator_agent_runtime_accepts_steer(&parent.state)
|| parent.state.status == "pending"
{
(
parent.state.run_id.as_str(),
parent.state.status.as_str(),
parent.state.phase.as_str(),
)
} else if let Some(task) = pending {
(
task.run_id.as_str(),
task.status.as_str(),
task.phase.as_str(),
)
} else {
writeln!(output, "[恢复] 当前 Session 没有可恢复的运行任务。")
.map_err(|error| format!("写入终端失败:{error}"))?;
return Ok(SwarmChatFlow::Continue);
};
writeln!(
output,
"[恢复] 继续观察 run={} status={} phase={}",
parent.state.run_id, parent.state.status, parent.state.phase
target_run_id, target_status, target_phase
)
.map_err(|error| format!("写入终端失败:{error}"))?;
let mut conversation_baseline =
new_swarm_turn_conversation_baseline(previous_message_count, &parent.state.run_id);
new_swarm_turn_conversation_baseline(previous_message_count, target_run_id);
capture_recovered_swarm_assistant_at(
root,
parent_agent_id,
@@ -27,11 +27,18 @@ pub(super) fn wait_for_swarm_turn<W: Write>(
let mut input_closed = false;
loop {
let runtimes = read_game_creator_agent_runtimes_at(root)?;
let changed = observer.print_changes(&runtimes, output)?;
let turn_runtimes = swarm_current_runtimes_for_run(
parent_agent_id,
session_id,
run_profile,
&conversation_baseline.parent_run_id,
&runtimes,
);
let changed = observer.print_changes(&turn_runtimes, output)?;
if changed {
stable_since = None;
}
let mut reconciliation = swarm_reconciliation_agents(&runtimes);
let mut reconciliation = swarm_reconciliation_agents(&turn_runtimes);
if !reconciliation.is_empty() {
observer.close_response_line(output)?;
return build_reconciliation_turn_outcome(
@@ -44,7 +51,13 @@ pub(super) fn wait_for_swarm_turn<W: Write>(
);
}
if !input_closed {
match observer.resolve_confirmations(root, parent_agent_id, &runtimes, input, output)? {
match observer.resolve_confirmations(
root,
parent_agent_id,
&turn_runtimes,
input,
output,
)? {
SwarmConfirmationResolution::Handled => {
stable_since = None;
recovery_scan_required = true;
@@ -61,7 +74,7 @@ pub(super) fn wait_for_swarm_turn<W: Write>(
match observer.resolve_user_input_requests(
root,
parent_agent_id,
&runtimes,
&turn_runtimes,
input,
output,
)? {
@@ -80,7 +93,7 @@ pub(super) fn wait_for_swarm_turn<W: Write>(
}
}
let pending_interactions =
swarm_unhandled_interaction_reasons(parent_agent_id, &runtimes, input_closed);
swarm_unhandled_interaction_reasons(parent_agent_id, &turn_runtimes, input_closed);
if !pending_interactions.is_empty() {
observer.close_response_line(output)?;
return build_incomplete_turn_outcome(
@@ -97,6 +110,7 @@ pub(super) fn wait_for_swarm_turn<W: Write>(
parent_agent_id,
session_id,
run_profile,
&conversation_baseline.parent_run_id,
&runtimes,
);
if !failure_scan.reconciliation_agents.is_empty() {
@@ -135,7 +149,15 @@ pub(super) fn wait_for_swarm_turn<W: Write>(
if last_runner_check.elapsed() >= Duration::from_secs(2) {
let runner = read_external_agent_runner_status();
last_runner_check = Instant::now();
if runtimes_are_busy(&runtimes) && (!runner.enabled || !runner.running) {
if swarm_turn_is_busy(
root,
parent_agent_id,
session_id,
run_profile,
&conversation_baseline.parent_run_id,
&runtimes,
)? && (!runner.enabled || !runner.running)
{
reconciliation.push("external-runner".to_string());
observer.close_response_line(output)?;
return build_reconciliation_turn_outcome(
@@ -148,7 +170,14 @@ pub(super) fn wait_for_swarm_turn<W: Write>(
);
}
}
if runtimes_are_busy(&runtimes) {
if swarm_turn_is_busy(
root,
parent_agent_id,
session_id,
run_profile,
&conversation_baseline.parent_run_id,
&runtimes,
)? {
stable_since = None;
recovery_scan_required = true;
} else {
@@ -175,13 +204,37 @@ pub(super) fn wait_for_swarm_turn<W: Write>(
session_id,
&conversation_baseline,
)?;
let parent_runtime =
swarm_parent_runtime(parent_agent_id, session_id, run_profile, &runtimes);
let completion_blockers = parent_runtime
.map(|parent| swarm_parent_completion_contract_blockers_at(root, parent))
.unwrap_or_else(|| vec!["parent-runtime-missing".to_string()]);
let parent_runtime_is_current = swarm_parent_runtime_for_run(
parent_agent_id,
session_id,
run_profile,
&conversation_baseline.parent_run_id,
&runtimes,
)
.is_some();
let parent_runtime = swarm_parent_runtime_snapshot_for_run(
root,
parent_agent_id,
session_id,
run_profile,
&conversation_baseline.parent_run_id,
&runtimes,
)?;
let completion_blockers = if parent_runtime_is_current {
parent_runtime
.as_ref()
.map(|parent| swarm_parent_completion_contract_blockers_at(root, parent))
.unwrap_or_else(|| vec!["parent-runtime-missing".to_string()])
} else if parent_runtime
.as_ref()
.is_some_and(parent_runtime_completed)
{
Vec::new()
} else {
vec!["parent-runtime-missing".to_string()]
};
match classify_swarm_turn_terminal(
parent_runtime,
parent_runtime.as_ref(),
conversation_metrics,
0,
0,
@@ -210,10 +263,11 @@ pub(super) fn wait_for_swarm_turn<W: Write>(
SwarmTurnReportOutcome::Settled,
parent_agent_id,
session_id,
Some(&conversation_baseline.parent_run_id),
&runtimes,
conversation_metrics,
0,
);
)?;
return Ok(SwarmTurnOutcome::Settled(report));
}
SwarmTurnTerminalClassification::Failed => {
@@ -227,6 +281,7 @@ pub(super) fn wait_for_swarm_turn<W: Write>(
"{}:{}",
parent_agent_id,
parent_runtime
.as_ref()
.map(|runtime| runtime.state.phase.as_str())
.unwrap_or("missing")
)],
@@ -236,7 +291,7 @@ pub(super) fn wait_for_swarm_turn<W: Write>(
let mut reasons = completion_blockers;
append_swarm_terminal_snapshot_reasons(
&mut reasons,
parent_runtime,
parent_runtime.as_ref(),
conversation_metrics,
);
return build_incomplete_turn_outcome(
@@ -284,12 +339,13 @@ pub(super) fn wait_for_swarm_turn<W: Write>(
.map_err(|error| format!("写入终端失败:{error}"))?;
}
SwarmChatInput::Message(message) => {
if let Some(parent) = runtimes.iter().find(|runtime| {
runtime.state.agent_id == parent_agent_id
&& runtime.state.session_id == session_id
&& runtime.state.run_profile == run_profile
&& matches!(runtime.state.status.as_str(), "pending" | "running")
}) {
if let Some(parent) = swarm_parent_steer_target(
parent_agent_id,
session_id,
run_profile,
Some(&conversation_baseline.parent_run_id),
&runtimes,
) {
let steer_id = format!("swarm-steer-{}", unix_millis());
let result = steer_game_creator_agent_runtime_task(
root.display().to_string(),
@@ -1195,7 +1195,7 @@ fn background_agent_runtime_reconciliation_without_ledger_still_blocks_recovery(
}
#[tokio::test]
async fn background_agent_runtime_recovers_pending_task() {
async fn background_agent_runtime_recovers_pending_task_after_cancelled_canonical_run() {
let root = unique_project_path();
init_local_game_project_at(&root, "project-1", "月光厨房").expect("project init");
let (sender, receiver) = mpsc::channel();
@@ -1219,6 +1219,49 @@ async fn background_agent_runtime_recovers_pending_task() {
}}
}}"#
));
let session_id = "agent-session-design-director";
let mut cancelled =
default_game_creator_agent_runtime_state("design-director", "design-cancelled-run");
cancelled.session_id = session_id.to_string();
cancelled.status = "cancelled".to_string();
cancelled.phase = "cancelled".to_string();
cancelled.current_task = "此前任务已取消".to_string();
cancelled.current_action = "此前任务已经取消".to_string();
write_game_creator_agent_runtime_state(&root, &cancelled)
.expect("persist cancelled canonical runtime");
write_agent_runtime_task_record_for_test(
&root,
&AgentRuntimeTaskRecord {
goal_id: None,
goal_revision: 0,
goal_status: None,
schema_version: AGENT_RUNTIME_SCHEMA_VERSION.to_string(),
agent_id: "design-director".to_string(),
task_id: "design-director".to_string(),
session_id: session_id.to_string(),
run_id: "design-cancelled-run".to_string(),
source: "agent-background-task".to_string(),
run_profile: default_agent_runtime_run_profile(),
run_profile_binding_fingerprint: String::new(),
parent_agent_id: None,
parent_run_id: None,
delegation_id: None,
task: "此前任务已取消".to_string(),
status: "cancelled".to_string(),
phase: "cancelled".to_string(),
current_action: "此前任务已经取消".to_string(),
terminal_detail: Some("此前任务取消".to_string()),
error: None,
updated_at: unix_timestamp().saturating_sub(120),
},
);
crate::write_game_creator_agent_runtime_cancel_request(
&root,
"design-director",
"design-cancelled-run",
"保留旧 run 取消 tombstone",
)
.expect("persist old cancellation tombstone");
let task_record = AgentRuntimeTaskRecord {
goal_id: None,
goal_revision: 0,
@@ -1226,7 +1269,7 @@ async fn background_agent_runtime_recovers_pending_task() {
schema_version: AGENT_RUNTIME_SCHEMA_VERSION.to_string(),
agent_id: "design-director".to_string(),
task_id: "design-director".to_string(),
session_id: "agent-session-design-director".to_string(),
session_id: session_id.to_string(),
run_id: "design-pending-recover-run".to_string(),
source: "agent-background-task".to_string(),
run_profile: default_agent_runtime_run_profile(),
@@ -1243,11 +1286,33 @@ async fn background_agent_runtime_recovers_pending_task() {
updated_at: unix_timestamp().saturating_sub(60),
};
write_agent_runtime_task_record_for_test(&root, &task_record);
crate::append_local_conversation_message_for_session_idempotent_at(
&root,
Some("design-director"),
Some(session_id),
LocalConversationMessage {
role: "user".to_string(),
content: "恢复排队后台任务".to_string(),
agent_id: None,
},
"runtime-pending-after-completed",
)
.expect("persist queued user message once");
let before = read_game_creator_agent_runtime_at(&root, "design-director")
.expect("read completed runtime with pending queue");
assert_eq!(before.state.run_id, "design-cancelled-run");
assert_eq!(before.state.status, "cancelled");
assert_eq!(before.state.phase, "cancelled");
assert_eq!(before.task_queue.pending, 1);
let resumed =
resume_game_creator_agent_background_tasks_at(&root).expect("resume pending task");
assert_eq!(resumed.len(), 1);
assert_eq!(resumed[0].state.run_id, "design-pending-recover-run");
assert_eq!(
resumed[0].state.run_id, "design-pending-recover-run",
"unexpected recovery snapshot: {resumed:?}"
);
assert_eq!(resumed[0].state.status, "running");
assert_eq!(
resumed[0].state.current_action,
@@ -1266,6 +1331,20 @@ async fn background_agent_runtime_recovers_pending_task() {
);
let agent_db = fs::read_to_string(root.join(".agent/agent.db")).expect("agent db");
assert!(agent_db.contains("\"recoveredFromStatus\":\"pending\""));
let conversation =
read_local_conversation_for_session_at(&root, Some("design-director"), Some(session_id))
.expect("read recovered conversation");
assert_eq!(
conversation
.messages
.iter()
.filter(|message| {
message.role == "user" && message.content == "恢复排队后台任务"
})
.count(),
1,
"recovery must not duplicate the already-persisted user turn"
);
fs::remove_dir_all(root).ok();
}
@@ -5301,3 +5301,4 @@
- 并发边界:External Editor 项目、素材库、生图和下载请求都在项目锁外执行,请求阶段只能读取现有 manifest 快照,不能补 seed task 或重写 manifest;下载完成后先按声明的图片类型校验 PNG/JPEG/WebP/GIF 魔数,再取得项目写锁,复核协作边界与输出路径并提交文件、manifest、revision 和验证凭证。专业 Agent 的完成门禁要求同一 run 的 `verifiedRevision >= mutationRevision`,不能被并行 Agent 的后续全局 revision 误判为过期,也允许更晚 revision 的复验覆盖本人修改;当前全局 revision 的集成验证仍由 Project Supervisor 精确负责。
- 复用边界:已有规范素材只有同时满足固定路径、`art-spritesheet / image/* / canvas` 登记、本地普通文件存在且 PNG 签名有效时才跳过生成;空文件、伪 PNG 或损坏占位必须重新进入美术委派。
- 恢复边界:首批 policy snapshot、durable provider batch 与恢复校验继续绑定同一 required Agent 集合;格式修复必须根据当前 policy 补齐可选的 `art-asset-plan` 固定产物合同,不能只修复程序与质量委派后绕过美术交付。
- 2026-07-25 决策:Swarm CLI 的 `busy` 仅表示当前 Agent 仍有运行、等待、reconciliation 或排队工作,不能作为 steer 目标判定。steer 能力统一复用 Runtime 协议层门禁;terminal canonical run 即使汇总出 pending queue 也永远不可 steer。CLI 发现旧 completed/cancelled run 后仍有 pending run 时先通知独立 Runner 恢复;旧 cancelled run 的 tombstone 不能让恢复扫描跳过后续 pending。新 turn 以 mutation 返回的 `acceptedRunId` 为权威 baseline,失败、交互、收束和 `turn.report.parentRunId` 都只归属该 runcanonical 已推进到下一 run 时从 append-only task journal 恢复目标 run 终态。相同且已落盘的最后一条用户消息只恢复观察原 run,不能重复写 conversation、steer ledger 或 task ledger。active Goal 也必须精确匹配当前 Runtime 身份与可 steer 状态,不能只依据 Goal 的 `active` 字符串直接追加。连续 run 的 assistant 回复按 `agentId + sessionId + runId` 派生的 finalization message ID 归属,失败终态按 `(agentId, runId)` 聚合且完整 task journal 优先于可能滞后的 state 投影;`turn.report` 的运行和队列计数同样读取完整 journal 并仅统计目标 parent run 及其直接 children。
@@ -3611,3 +3611,11 @@
- 处理:正式 AppData 只作只读配置来源。每次人工测试在系统临时根创建 `0700` sentinel 隔离目录,只把主配置和可选 local overlay 私有复制为 `0600` 普通文件;不得复制 endpoint、lock、`.previous` 或其它状态。LLM 检查与 Swarm CLI 全部使用隔离目录。退出时通过内部 CLI 请求 `runner.shutdown_if_idle`,确认隔离 endpoint 消失后才删除配置;仍有任务或无法确认退出时同时保留测试项目和隔离配置并报告路径。正式 Runner 的 PID、bootId、端口和 executable fingerprint 必须保持不变。
- 验证:单元测试覆盖私有 inode、权限、local overlay、禁止复制 endpoint/lock/备份、符号链接拒绝、sentinel 清理和 endpoint 存在时拒绝删除;真实 smoke 使用隔离 AppData 启动并收束空闲 Runner,前后比较正式 endpoint 身份且确认正式 PID 存活,再检查本轮 `/tmp` 项目和隔离配置均已清理。
- 关联:`apps/ai-game-creator-shell/scripts/agent-swarm-test-chat.mjs``apps/ai-game-creator-shell/tests/agentSwarmTestEntry.test.ts``apps/ai-game-creator-shell/src-tauri/src/runner/client.rs``apps/ai-game-creator-shell/src-tauri/src/cli.rs`
## Swarm 队列 busy 不能直接当成 canonical run 可 steer
- 现象:继续已有项目时,Runtime state 仍指向旧的 `idle / completed``cancelled / cancelled` run A,但 task ledger 已有更新的 `pending / queued` run B;CLI 打印“已投递 B”后却立刻把 A 及其历史子 Agent 的 cancelled/budget-exhausted 报成 B 的失败。
- 原因:旧 `runtime_is_busy` 同时包含当前 state 和队列汇总,调用方看到 `task_queue.pending > 0` 后仍从 canonical state 反推 steer、失败扫描和 turn report 的 runId;取消 tombstone 还会让恢复扫描在处理 A 后无条件跳过 B。底层拒绝 terminal steer 和保留 A 的真实失败历史都是正确行为,不能通过放宽门禁或删除历史记录修复。
- 处理:保留 queue busy 用于 Runner 存活判断,另由 Runtime 协议层提供唯一 steerable 判定。start mutation 返回实际 `acceptedRunId`CLI 以它建立不可变 turn baseline;失败、reconciliation、用户交互、收束和报告只观察该 run。canonical 已推进到后续 run 时从 task journal 读取目标 run 的最终记录。旧 cancelled canonical 若仍有 pending 且无 running,恢复扫描跳过旧 run 的 pending action 恢复,直接启动队首 pending。若输入与已落盘 pending task 及最后一条 user 消息相同,则只观察原 run。Goal 路径也必须核对同一 Agent、Session、runId、Run Profile 和 steerable 状态。连续 run 的回复必须按确定性 finalization message ID 过滤;历史 specialist 失败必须以 `(agentId, runId)` 为键读取完整 journal,不能让滞后的非失败 state 删除 journal 已记录的失败;报告计数也不能退回 `recent_tasks` 的 12 条窗口。
- 验证:构造 cancelled run A、保留 A cancel tombstone、pending run B 和单份已落盘用户消息,证明恢复后 B 进入 running 并完成且 conversation 不重复。另覆盖观察 B 时忽略 A 及 A 子任务失败、观察 A 时仍正常失败、B 完成后 canonical 已推进到 C 仍可从 journal 收束 B、`turn.report.parentRunId` 始终为 baseline,以及 expected Goal runId 不一致时不选中目标。
- 关联:`apps/ai-game-creator-shell/src-tauri/src/swarm_cli/turn_dispatch.rs``apps/ai-game-creator-shell/src-tauri/src/swarm_cli/terminal_classification.rs``apps/ai-game-creator-shell/src-tauri/src/agent/runtime_protocol/steering.rs``apps/ai-game-creator-shell/src-tauri/src/agent/runtime_driver/recovery_scan.rs`
@@ -717,3 +717,4 @@ game-project/
- 2026-07-22 的 V1.46 消除自主构建的人工确认等待。`autonomous-game-build` 仅把固定 auto-safe 白名单提升为自动执行;项目级、Agent 级和 MCP catalog 动态策略产生的其余 `RequiresConfirmation` 一律转为 `Denied`,向同一 run 返回“改用 auto-safe 工具或省略动作”的 observation。显式 deny 继续优先,标准 profile 的确认语义不变;包含拒绝成员的 Provider action 批次在任何工具执行前整体 `aborted`,不能执行 auto 前缀。
- V1.46 的 Runner 恢复会把旧版本已持久化的自主 `pending-confirmation / waiting-confirmation` 迁移为 `observed-rejected / aborted`。Provider batch 是原子提交点,pending 只是可重建镜像;批次先落盘、pending 后落盘之间强杀时,下次恢复从 aborted batch 补齐拒绝账本并继续原 Session/run。公共状态和审计使用 `runtime-policy-rejected`,不得写成“已执行自动工具”或“开发者拒绝”。
- V1.46 稳定树通过自主构建过滤 `20/20`、确认相关 `18/18`、旧等待态完整恢复、既有批次零副作用、`cargo check --tests`、Rust 串行全量 `1149 passed / 5 ignored / 0 failed`、fmt 与 diff 检查。确定性正式 E2E 为 **PASS**Provider lifecycle `17/17`、revision `0 -> 2`、Chrome `37/37`、人工输入与残留均为 `0`。新的独立外部 Provider 同轮 E2E 也为 **PASS**:单条任务后 EOFapprove / answer / steer 为 `0`Provider lifecycle `62/62`revision `0 -> 5``game/index.html``7816` bytesstatic smoke、desktop / mobile 与 `lane-defense-v1 37/37` 全通过,唯一 Supervisor assistant,全部 sidecar、reconciliation、重复和泄漏计数均为 `0`
- 2026-07-25 补充:Swarm CLI 将“Agent 执行通道仍忙”和“canonical run 可接受 steer”拆为两个判定。`task_queue.pending > 0` 继续用于 Runner 存活观察,但只有底层 steer 门禁认可的 `running / waiting-for-confirmation` Runtime 才能接收新指令;`completed/cancelled run A + pending run B` 必须先通过 Runner 恢复 B,禁止把消息追加到 A。start mutation 必须返回实际 `acceptedRunId`CLI 以它建立 turn baseline,失败扫描、交互、收束和 `turn.report` 都只认 baseline run;若 Runner 已连续推进到 C,则从 task journal 读取 B 的终态。旧 A 的 cancel tombstone 不得阻断队首 B 恢复。若用户重复输入的内容正是已经落盘的 pending 任务且它仍是对话最后一条 user 消息,CLI 只观察原 run,不再追加第二份对话或创建新任务;active Goal 同样必须同时匹配 Agent、Session、runId、Run Profile 和可 steer 状态。连续 B/C 的 assistant 回复按 finalization message ID 归属,observer 不输出非目标 run 的状态或流;历史失败和报告计数都读取完整 task journal,失败聚合使用 `(agentId, runId)`,journal 失败终态不能被滞后的 state 投影覆盖。