完成素材无限画布最终审计与恢复安全收口
同步资源管理最新稳定的 manifest relay、Runner 重挂与 MCP 域名修复 统一网站与 Tauri 的共享画布 history、缩放恢复和宿主边界 移除生成账本中的临时上传凭证与 Provider 敏感字段 补齐事务快照发布前故障恢复和同幂等身份重放 更新最终 PRD、技术方案、decision-log 与 pitfalls 补齐前端、Tauri Rust、竞态、锁、恢复和原生壳门禁
This commit is contained in:
@@ -5,6 +5,9 @@ pub(super) static GAME_CREATOR_AGENT_RUNTIME_UPDATE_APP_HANDLE: OnceLock<tauri::
|
||||
pub(super) static GAME_CREATOR_MANIFEST_INVALIDATION_EVENT_SINK: OnceLock<
|
||||
std::sync::Mutex<Option<GameCreatorManifestInvalidationEventSink>>,
|
||||
> = OnceLock::new();
|
||||
#[cfg(test)]
|
||||
pub(super) static GAME_CREATOR_MANIFEST_INVALIDATION_EVENT_SINK_TEST_LOCK: std::sync::Mutex<()> =
|
||||
std::sync::Mutex::new(());
|
||||
pub(super) static STATIC_DELEGATE_PARENT_WAKE_SINGLEFLIGHT: OnceLock<
|
||||
std::sync::Mutex<std::collections::BTreeMap<String, bool>>,
|
||||
> = OnceLock::new();
|
||||
@@ -238,7 +241,7 @@ pub(in crate::agent) use task_queue::*;
|
||||
pub(in crate::agent) use task_start::*;
|
||||
|
||||
#[cfg(test)]
|
||||
pub(crate) use entrypoints::clear_game_creator_manifest_invalidation_event_sink_for_test;
|
||||
pub(crate) use entrypoints::acquire_game_creator_manifest_invalidation_event_sink_test_guard;
|
||||
#[allow(unused_imports)]
|
||||
pub(crate) use entrypoints::{
|
||||
chat_with_game_creator_agent_at, chat_with_game_creator_role_agent_at,
|
||||
|
||||
@@ -80,8 +80,37 @@ pub(crate) fn configure_game_creator_manifest_invalidation_event_sink(
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
pub(crate) fn clear_game_creator_manifest_invalidation_event_sink_for_test() {
|
||||
*lock_game_creator_manifest_invalidation_event_sink() = None;
|
||||
pub(crate) struct GameCreatorManifestInvalidationEventSinkTestGuard {
|
||||
_isolation: std::sync::MutexGuard<'static, ()>,
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
impl GameCreatorManifestInvalidationEventSinkTestGuard {
|
||||
pub(crate) fn configure(&self, port: u16, token: &str) -> Result<(), String> {
|
||||
configure_game_creator_manifest_invalidation_event_sink(port, token)
|
||||
}
|
||||
|
||||
pub(crate) fn configured_sink(&self) -> Option<GameCreatorManifestInvalidationEventSink> {
|
||||
lock_game_creator_manifest_invalidation_event_sink().clone()
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
impl Drop for GameCreatorManifestInvalidationEventSinkTestGuard {
|
||||
fn drop(&mut self) {
|
||||
*lock_game_creator_manifest_invalidation_event_sink() = None;
|
||||
}
|
||||
}
|
||||
|
||||
#[cfg(test)]
|
||||
pub(crate) fn acquire_game_creator_manifest_invalidation_event_sink_test_guard(
|
||||
) -> GameCreatorManifestInvalidationEventSinkTestGuard {
|
||||
let isolation = GAME_CREATOR_MANIFEST_INVALIDATION_EVENT_SINK_TEST_LOCK
|
||||
.lock()
|
||||
.unwrap_or_else(|poisoned| poisoned.into_inner());
|
||||
GameCreatorManifestInvalidationEventSinkTestGuard {
|
||||
_isolation: isolation,
|
||||
}
|
||||
}
|
||||
|
||||
fn relay_game_creator_manifest_invalidation(root: &Path, agent_id: &str) -> Result<(), String> {
|
||||
|
||||
@@ -2229,6 +2229,9 @@ fn validate_generation_provenance(value: &AssetCanvasGenerationProvenance) -> Re
|
||||
#[cfg(test)]
|
||||
#[derive(Clone, Copy, Debug, Eq, PartialEq)]
|
||||
enum AssetCanvasCommitFaultStage {
|
||||
FirstSnapshotInstalled,
|
||||
SnapshotsInstalled,
|
||||
JournalInstalled,
|
||||
Prepared,
|
||||
FileInstalled,
|
||||
ManifestInstalled,
|
||||
@@ -2248,7 +2251,16 @@ fn maybe_fail_asset_canvas_commit(
|
||||
if fault.is_some_and(|fault| {
|
||||
matches!(
|
||||
(fault, stage),
|
||||
(AssetCanvasCommitFaultStage::Prepared, "prepared")
|
||||
(
|
||||
AssetCanvasCommitFaultStage::FirstSnapshotInstalled,
|
||||
"first-snapshot-installed"
|
||||
) | (
|
||||
AssetCanvasCommitFaultStage::SnapshotsInstalled,
|
||||
"snapshots-installed"
|
||||
) | (
|
||||
AssetCanvasCommitFaultStage::JournalInstalled,
|
||||
"journal-installed"
|
||||
) | (AssetCanvasCommitFaultStage::Prepared, "prepared")
|
||||
| (AssetCanvasCommitFaultStage::FileInstalled, "file-installed")
|
||||
| (
|
||||
AssetCanvasCommitFaultStage::ManifestInstalled,
|
||||
@@ -2590,6 +2602,7 @@ fn commit_asset_canvas_at_internal(
|
||||
"manifest.before.json",
|
||||
&manifest_before_bytes,
|
||||
)?;
|
||||
maybe_fail_asset_canvas_commit(fault, "first-snapshot-installed")?;
|
||||
write_asset_canvas_transaction_snapshot(
|
||||
root,
|
||||
&input.commit_id,
|
||||
@@ -2608,7 +2621,9 @@ fn commit_asset_canvas_at_internal(
|
||||
"project-revision.after.json",
|
||||
&revision_after_bytes,
|
||||
)?;
|
||||
maybe_fail_asset_canvas_commit(fault, "snapshots-installed")?;
|
||||
write_asset_canvas_journal(root, &journal)?;
|
||||
maybe_fail_asset_canvas_commit(fault, "journal-installed")?;
|
||||
write_asset_canvas_ledger(root, &ledger)?;
|
||||
maybe_fail_asset_canvas_commit(fault, "prepared")?;
|
||||
|
||||
@@ -3079,6 +3094,92 @@ fn recover_asset_canvas_transaction_locked(
|
||||
mark_asset_canvas_reconciliation_locked(root, journal, ledger).map(|outcome| (outcome, None))
|
||||
}
|
||||
|
||||
fn clean_unpublished_asset_canvas_transaction_locked(
|
||||
root: &Path,
|
||||
commit_id: &str,
|
||||
journal: Option<&AssetCanvasTransactionJournal>,
|
||||
) -> Result<RecoverAssetCanvasOutcome, String> {
|
||||
if let Some(journal) = journal {
|
||||
if journal.schema_version != ASSET_CANVAS_TRANSACTION_SCHEMA_VERSION
|
||||
|| journal.commit_id != commit_id
|
||||
|| journal.stage != AssetCanvasTransactionStage::Prepared
|
||||
{
|
||||
return Err("缺少 ledger 的素材画布 transaction 身份无效".to_string());
|
||||
}
|
||||
let final_path = resolve_local_project_path(root, &journal.final_image_relative_path)?;
|
||||
match fs::symlink_metadata(&final_path) {
|
||||
Err(error) if error.kind() == std::io::ErrorKind::NotFound => {}
|
||||
Ok(_) => return Err("缺少 ledger 的素材画布 transaction 已产生正式文件".to_string()),
|
||||
Err(_) => return Err("检查未发布素材画布 transaction 正式文件失败".to_string()),
|
||||
}
|
||||
let current_manifest = current_asset_canvas_manifest(root)?;
|
||||
if asset_canvas_sha256(&asset_canvas_json_bytes(¤t_manifest)?)
|
||||
!= journal.manifest_before_sha256
|
||||
{
|
||||
return Err("缺少 ledger 的素材画布 transaction manifest 已变化".to_string());
|
||||
}
|
||||
let current_revision = read_game_creator_agent_runtime_project_revision(root)?;
|
||||
if journal.project_revision_before_sha256.as_deref()
|
||||
!= Some(asset_canvas_sha256(&asset_canvas_json_bytes(¤t_revision)?).as_str())
|
||||
{
|
||||
return Err("缺少 ledger 的素材画布 transaction revision 已变化".to_string());
|
||||
}
|
||||
}
|
||||
|
||||
let transaction_directory = resolve_local_project_path(
|
||||
root,
|
||||
&format!("{ASSET_CANVAS_ROOT}/transactions/{commit_id}"),
|
||||
)?;
|
||||
let mut files = Vec::new();
|
||||
for entry in fs::read_dir(&transaction_directory)
|
||||
.map_err(|_| "读取未发布素材画布 transaction 失败".to_string())?
|
||||
{
|
||||
let entry = entry.map_err(|_| "读取未发布素材画布 transaction 条目失败".to_string())?;
|
||||
let metadata = fs::symlink_metadata(entry.path())
|
||||
.map_err(|_| "读取未发布素材画布 transaction 元数据失败".to_string())?;
|
||||
let name = entry.file_name().to_string_lossy().into_owned();
|
||||
let is_expected_snapshot = matches!(
|
||||
name.as_str(),
|
||||
"manifest.before.json"
|
||||
| "manifest.after.json"
|
||||
| "project-revision.before.json"
|
||||
| "project-revision.after.json"
|
||||
| "journal.json"
|
||||
);
|
||||
let is_owned_temporary_file = (name.starts_with(".asset-canvas-")
|
||||
&& name.ends_with(".tmp"))
|
||||
|| name.starts_with(".journal.json.tmp.");
|
||||
if metadata.file_type().is_symlink()
|
||||
|| !metadata.is_file()
|
||||
|| metadata.len() > ASSET_CANVAS_MAX_LEDGER_BYTES as u64
|
||||
|| (!is_expected_snapshot && !is_owned_temporary_file)
|
||||
{
|
||||
return Err("未发布素材画布 transaction 包含未知文件".to_string());
|
||||
}
|
||||
files.push(entry.path());
|
||||
}
|
||||
if files.len() > 16 {
|
||||
return Err("未发布素材画布 transaction 文件数量超限".to_string());
|
||||
}
|
||||
for path in files {
|
||||
fs::remove_file(path).map_err(|_| "清理未发布素材画布 transaction 文件失败".to_string())?;
|
||||
}
|
||||
fs::remove_dir(&transaction_directory)
|
||||
.map_err(|_| "清理未发布素材画布 transaction 目录失败".to_string())?;
|
||||
#[cfg(unix)]
|
||||
if let Some(parent) = transaction_directory.parent() {
|
||||
File::open(parent)
|
||||
.and_then(|directory| directory.sync_all())
|
||||
.map_err(|_| "同步素材画布 transaction 清理结果失败".to_string())?;
|
||||
}
|
||||
Ok(RecoverAssetCanvasOutcome {
|
||||
commit_id: commit_id.to_string(),
|
||||
status: RecoverAssetCanvasOutcomeStatus::RolledBack,
|
||||
event_id: None,
|
||||
asset_id: None,
|
||||
})
|
||||
}
|
||||
|
||||
pub(crate) fn recover_asset_canvas_transactions_at(
|
||||
root: &Path,
|
||||
expected_project_id: &str,
|
||||
@@ -3114,7 +3215,16 @@ pub(crate) fn recover_asset_canvas_transactions_at(
|
||||
let mut outcomes = Vec::new();
|
||||
let mut events = Vec::new();
|
||||
for commit_id in commit_ids {
|
||||
let Some(journal) = read_asset_canvas_journal(root, &commit_id)? else {
|
||||
let journal = read_asset_canvas_journal(root, &commit_id)?;
|
||||
if read_asset_canvas_ledger(root, &commit_id)?.is_none() {
|
||||
outcomes.push(clean_unpublished_asset_canvas_transaction_locked(
|
||||
root,
|
||||
&commit_id,
|
||||
journal.as_ref(),
|
||||
)?);
|
||||
continue;
|
||||
}
|
||||
let Some(journal) = journal else {
|
||||
return Err("素材画布 transaction 缺少 journal".to_string());
|
||||
};
|
||||
let (outcome, event) = recover_asset_canvas_transaction_locked(root, journal)?;
|
||||
|
||||
@@ -128,14 +128,12 @@ impl From<ExternalCanvasGenerationContext> for PrivateCanvasContext {
|
||||
}
|
||||
}
|
||||
|
||||
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
|
||||
#[serde(deny_unknown_fields, rename_all = "camelCase")]
|
||||
#[derive(Clone)]
|
||||
struct PrivateUploadTicket {
|
||||
host: String,
|
||||
bucket: String,
|
||||
object_key: String,
|
||||
success_action_status: u16,
|
||||
max_size_bytes: u64,
|
||||
form_fields: BTreeMap<String, String>,
|
||||
}
|
||||
|
||||
@@ -145,7 +143,8 @@ struct PrivateReferenceState {
|
||||
resource_id: String,
|
||||
stable_reference: Option<String>,
|
||||
asset_object_id: Option<String>,
|
||||
upload_ticket: Option<PrivateUploadTicket>,
|
||||
upload_bucket: Option<String>,
|
||||
upload_object_key: Option<String>,
|
||||
upload_completed: bool,
|
||||
}
|
||||
|
||||
@@ -336,6 +335,21 @@ fn validate_generation_ledger(ledger: &AssetCanvasGenerationLedger) -> Result<()
|
||||
{
|
||||
return Err("素材画布私有生成账本内容无效".to_string());
|
||||
}
|
||||
for state in &ledger.reference_states {
|
||||
if state.resource_id.trim().is_empty()
|
||||
|| state.resource_id.chars().count() > 512
|
||||
|| state.resource_id.chars().any(char::is_control)
|
||||
|| state.upload_bucket.is_some() != state.upload_object_key.is_some()
|
||||
|| state.upload_bucket.as_ref().is_some_and(|value| {
|
||||
value.is_empty() || value.len() > 512 || value.chars().any(char::is_control)
|
||||
})
|
||||
|| state.upload_object_key.as_ref().is_some_and(|value| {
|
||||
value.is_empty() || value.len() > 2048 || value.chars().any(char::is_control)
|
||||
})
|
||||
{
|
||||
return Err("素材画布私有生成参考状态无效".to_string());
|
||||
}
|
||||
}
|
||||
if let (Some(body), Some(expected_sha)) = (
|
||||
ledger.request_body_json.as_deref(),
|
||||
ledger.request_body_sha256.as_deref(),
|
||||
@@ -756,7 +770,8 @@ async fn try_confirm_uploaded_reference(
|
||||
client: &reqwest::Client,
|
||||
api_base_url: &str,
|
||||
api_key: &str,
|
||||
ticket: &PrivateUploadTicket,
|
||||
bucket: &str,
|
||||
object_key: &str,
|
||||
material: &ReferenceMaterial,
|
||||
asset_kind: &str,
|
||||
) -> Result<Option<String>, String> {
|
||||
@@ -766,8 +781,8 @@ async fn try_confirm_uploaded_reference(
|
||||
))
|
||||
.bearer_auth(api_key)
|
||||
.json(&serde_json::json!({
|
||||
"bucket": ticket.bucket,
|
||||
"objectKey": ticket.object_key,
|
||||
"bucket": bucket,
|
||||
"objectKey": object_key,
|
||||
"contentType": material.media_type,
|
||||
"contentLength": material.bytes.as_ref().map(Vec::len),
|
||||
"contentHash": material.sha256,
|
||||
@@ -790,7 +805,7 @@ async fn try_confirm_uploaded_reference(
|
||||
.unwrap_or(&serde_json::Value::Null);
|
||||
let confirmed_object_key = json_string_field(asset_object, "objectKey")
|
||||
.ok_or_else(|| "参考资源确认响应缺少 objectKey".to_string())?;
|
||||
if confirmed_object_key != ticket.object_key {
|
||||
if confirmed_object_key != object_key {
|
||||
return Err("参考资源确认响应 objectKey 不一致".to_string());
|
||||
}
|
||||
json_string_field(asset_object, "assetObjectId")
|
||||
@@ -867,7 +882,6 @@ async fn request_upload_ticket(
|
||||
bucket,
|
||||
object_key,
|
||||
success_action_status,
|
||||
max_size_bytes,
|
||||
form_fields,
|
||||
})
|
||||
}
|
||||
@@ -941,7 +955,8 @@ async fn ensure_reference_states(
|
||||
resource_id: resource_id.clone(),
|
||||
stable_reference: material.stable_reference.clone(),
|
||||
asset_object_id: None,
|
||||
upload_ticket: None,
|
||||
upload_bucket: None,
|
||||
upload_object_key: None,
|
||||
upload_completed: material.stable_reference.is_some(),
|
||||
});
|
||||
write_generation_ledger(root, ledger)?;
|
||||
@@ -957,19 +972,25 @@ async fn ensure_reference_states(
|
||||
write_generation_ledger(root, ledger)?;
|
||||
continue;
|
||||
}
|
||||
if let Some(ticket) = ledger.reference_states[index].upload_ticket.clone() {
|
||||
if let (Some(bucket), Some(object_key)) = (
|
||||
ledger.reference_states[index].upload_bucket.clone(),
|
||||
ledger.reference_states[index].upload_object_key.clone(),
|
||||
) {
|
||||
if let Some(asset_object_id) = try_confirm_uploaded_reference(
|
||||
client,
|
||||
api_base_url,
|
||||
api_key,
|
||||
&ticket,
|
||||
&bucket,
|
||||
&object_key,
|
||||
&material,
|
||||
&ledger.asset_kind,
|
||||
)
|
||||
.await?
|
||||
{
|
||||
ledger.reference_states[index].stable_reference = Some(ticket.object_key);
|
||||
ledger.reference_states[index].stable_reference = Some(object_key);
|
||||
ledger.reference_states[index].asset_object_id = Some(asset_object_id);
|
||||
ledger.reference_states[index].upload_bucket = None;
|
||||
ledger.reference_states[index].upload_object_key = None;
|
||||
ledger.reference_states[index].upload_completed = true;
|
||||
write_generation_ledger(root, ledger)?;
|
||||
continue;
|
||||
@@ -977,7 +998,8 @@ async fn ensure_reference_states(
|
||||
}
|
||||
let ticket =
|
||||
request_upload_ticket(client, api_base_url, api_key, ledger, &material).await?;
|
||||
ledger.reference_states[index].upload_ticket = Some(ticket.clone());
|
||||
ledger.reference_states[index].upload_bucket = Some(ticket.bucket.clone());
|
||||
ledger.reference_states[index].upload_object_key = Some(ticket.object_key.clone());
|
||||
ledger.reference_states[index].upload_completed = false;
|
||||
write_generation_ledger(root, ledger)?;
|
||||
upload_reference(&ticket, &material, api_base_url).await?;
|
||||
@@ -987,7 +1009,8 @@ async fn ensure_reference_states(
|
||||
client,
|
||||
api_base_url,
|
||||
api_key,
|
||||
&ticket,
|
||||
&ticket.bucket,
|
||||
&ticket.object_key,
|
||||
&material,
|
||||
&ledger.asset_kind,
|
||||
)
|
||||
@@ -995,6 +1018,8 @@ async fn ensure_reference_states(
|
||||
.ok_or_else(|| "参考资源上传后无法确认稳定对象".to_string())?;
|
||||
ledger.reference_states[index].stable_reference = Some(ticket.object_key);
|
||||
ledger.reference_states[index].asset_object_id = Some(asset_object_id);
|
||||
ledger.reference_states[index].upload_bucket = None;
|
||||
ledger.reference_states[index].upload_object_key = None;
|
||||
write_generation_ledger(root, ledger)?;
|
||||
}
|
||||
ledger.resolved_reference_ids = ledger
|
||||
@@ -2283,14 +2308,16 @@ mod tests {
|
||||
resource_id: source_resource_id.clone(),
|
||||
stable_reference: Some("objects/source.png".to_string()),
|
||||
asset_object_id: Some("source-object".to_string()),
|
||||
upload_ticket: None,
|
||||
upload_bucket: None,
|
||||
upload_object_key: None,
|
||||
upload_completed: true,
|
||||
},
|
||||
PrivateReferenceState {
|
||||
resource_id: "local-asset:style".to_string(),
|
||||
stable_reference: Some("objects/style.png".to_string()),
|
||||
asset_object_id: Some("style-object".to_string()),
|
||||
upload_ticket: None,
|
||||
upload_bucket: None,
|
||||
upload_object_key: None,
|
||||
upload_completed: true,
|
||||
},
|
||||
],
|
||||
@@ -2334,6 +2361,64 @@ mod tests {
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn private_generation_ledger_never_serializes_upload_credentials_or_provider_url() {
|
||||
let project_id = "phase-five-private-upload-ledger";
|
||||
let (directory, draft) = create_generation_fixture(project_id, "阶段五私有上传账本测试");
|
||||
let ticket = PrivateUploadTicket {
|
||||
host: "https://private-upload.provider.example.test/signed".to_string(),
|
||||
bucket: "stable-private-bucket".to_string(),
|
||||
object_key: "asset-canvas-references/project/reference.png".to_string(),
|
||||
success_action_status: 204,
|
||||
form_fields: BTreeMap::from([
|
||||
(
|
||||
"Authorization".to_string(),
|
||||
"private-authorization".to_string(),
|
||||
),
|
||||
("policy".to_string(), "private-upload-policy".to_string()),
|
||||
(
|
||||
"signature".to_string(),
|
||||
"private-upload-signature".to_string(),
|
||||
),
|
||||
]),
|
||||
};
|
||||
let mut ledger = accepted_ledger(
|
||||
project_id,
|
||||
&draft,
|
||||
"https://editor.example.test",
|
||||
"private-api-key",
|
||||
);
|
||||
ledger.phase = GenerationLedgerPhase::ReferencesPreparing;
|
||||
ledger.reference_states = vec![PrivateReferenceState {
|
||||
resource_id: "local-asset:reference".to_string(),
|
||||
stable_reference: None,
|
||||
asset_object_id: None,
|
||||
upload_bucket: Some(ticket.bucket.clone()),
|
||||
upload_object_key: Some(ticket.object_key.clone()),
|
||||
upload_completed: false,
|
||||
}];
|
||||
ledger.requested_reference_resource_ids = vec!["local-asset:reference".to_string()];
|
||||
|
||||
let persisted = serde_json::to_string_pretty(&ledger).expect("serialize private ledger");
|
||||
assert!(persisted.contains("stable-private-bucket"));
|
||||
assert!(persisted.contains("asset-canvas-references/project/reference.png"));
|
||||
for forbidden in [
|
||||
"uploadTicket",
|
||||
"formFields",
|
||||
ticket.host.as_str(),
|
||||
"private-authorization",
|
||||
"private-upload-policy",
|
||||
"private-upload-signature",
|
||||
"private-api-key",
|
||||
] {
|
||||
assert!(
|
||||
!persisted.contains(forbidden),
|
||||
"private ledger leaked forbidden upload material: {forbidden}"
|
||||
);
|
||||
}
|
||||
drop(directory);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn unknown_poll_result_keeps_the_original_operation_in_get_only_recovery() {
|
||||
let project_id = "phase-five-poll-reconciliation";
|
||||
|
||||
@@ -544,6 +544,9 @@ fn same_project_revision_double_commit_cannot_overwrite() {
|
||||
#[test]
|
||||
fn every_commit_fault_stage_recovers_without_ambiguous_overwrite() {
|
||||
for fault in [
|
||||
AssetCanvasCommitFaultStage::FirstSnapshotInstalled,
|
||||
AssetCanvasCommitFaultStage::SnapshotsInstalled,
|
||||
AssetCanvasCommitFaultStage::JournalInstalled,
|
||||
AssetCanvasCommitFaultStage::Prepared,
|
||||
AssetCanvasCommitFaultStage::FileInstalled,
|
||||
AssetCanvasCommitFaultStage::ManifestInstalled,
|
||||
@@ -573,7 +576,11 @@ fn every_commit_fault_stage_recovers_without_ambiguous_overwrite() {
|
||||
.find(|outcome| outcome.commit_id == input.commit_id)
|
||||
.expect("recovery outcome");
|
||||
match fault {
|
||||
AssetCanvasCommitFaultStage::Prepared | AssetCanvasCommitFaultStage::FileInstalled => {
|
||||
AssetCanvasCommitFaultStage::FirstSnapshotInstalled
|
||||
| AssetCanvasCommitFaultStage::SnapshotsInstalled
|
||||
| AssetCanvasCommitFaultStage::JournalInstalled
|
||||
| AssetCanvasCommitFaultStage::Prepared
|
||||
| AssetCanvasCommitFaultStage::FileInstalled => {
|
||||
assert_eq!(outcome.status, RecoverAssetCanvasOutcomeStatus::RolledBack);
|
||||
assert!(!fixture
|
||||
.root()
|
||||
@@ -593,6 +600,20 @@ fn every_commit_fault_stage_recovers_without_ambiguous_overwrite() {
|
||||
assert_eq!(recovered_draft.status, AssetCanvasDraftStatus::Editing);
|
||||
assert_eq!(recovered_draft.revision, draft.revision);
|
||||
assert!(recovered_draft.pending_commit.is_none());
|
||||
if matches!(
|
||||
fault,
|
||||
AssetCanvasCommitFaultStage::FirstSnapshotInstalled
|
||||
| AssetCanvasCommitFaultStage::SnapshotsInstalled
|
||||
| AssetCanvasCommitFaultStage::JournalInstalled
|
||||
) {
|
||||
assert!(!fixture
|
||||
.root()
|
||||
.join(format!(
|
||||
".agent/workbench/asset-canvas/transactions/{}",
|
||||
input.commit_id
|
||||
))
|
||||
.exists());
|
||||
}
|
||||
}
|
||||
_ => {
|
||||
assert!(matches!(
|
||||
@@ -632,6 +653,60 @@ fn every_commit_fault_stage_recovers_without_ambiguous_overwrite() {
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn unpublished_snapshot_transaction_recovers_and_replays_the_same_idempotency_identity() {
|
||||
let fixture = initialize_fixture();
|
||||
let draft = add_imported_layer(&fixture, &fixture.draft);
|
||||
let staged = stage_image(&fixture, &draft);
|
||||
let input = commit_input(
|
||||
&fixture,
|
||||
&draft,
|
||||
&staged,
|
||||
Uuid::new_v4().to_string(),
|
||||
Uuid::new_v4().to_string(),
|
||||
);
|
||||
commit_asset_canvas_at_internal(
|
||||
fixture.root(),
|
||||
&input,
|
||||
Some(AssetCanvasCommitFaultStage::SnapshotsInstalled),
|
||||
)
|
||||
.expect_err("stop before journal publication");
|
||||
|
||||
let recovered = recover_asset_canvas_transactions_at(fixture.root(), PROJECT_ID)
|
||||
.expect("clean unpublished transaction");
|
||||
assert_eq!(
|
||||
recovered.result.outcomes,
|
||||
vec![RecoverAssetCanvasOutcome {
|
||||
commit_id: input.commit_id.clone(),
|
||||
status: RecoverAssetCanvasOutcomeStatus::RolledBack,
|
||||
event_id: None,
|
||||
asset_id: None,
|
||||
}]
|
||||
);
|
||||
|
||||
assert!(matches!(
|
||||
commit_asset_canvas_at(fixture.root(), &input)
|
||||
.expect("replay same idempotency identity after cleanup")
|
||||
.result,
|
||||
CommitAssetCanvasResult::Committed { .. }
|
||||
));
|
||||
assert!(matches!(
|
||||
commit_asset_canvas_at(fixture.root(), &input)
|
||||
.expect("repeat committed identity")
|
||||
.result,
|
||||
CommitAssetCanvasResult::AlreadyCommitted { .. }
|
||||
));
|
||||
assert_eq!(
|
||||
current_asset_canvas_manifest(fixture.root())
|
||||
.expect("manifest after replay")
|
||||
.assets
|
||||
.iter()
|
||||
.filter(|asset| asset.id.starts_with("canvas-"))
|
||||
.count(),
|
||||
1
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn mismatched_installed_file_requires_reconciliation_and_is_not_deleted() {
|
||||
let fixture = initialize_fixture();
|
||||
|
||||
@@ -11,6 +11,7 @@ use std::io::{self, BufRead, BufReader, Read, Write};
|
||||
use std::net::{Ipv4Addr, SocketAddrV4, TcpStream};
|
||||
use std::path::{Path, PathBuf};
|
||||
use std::process::{Child, Command, Stdio};
|
||||
use std::sync::{Mutex, OnceLock};
|
||||
use std::thread;
|
||||
use std::time::{Duration, Instant};
|
||||
|
||||
@@ -19,6 +20,76 @@ const AGENT_RUNNER_LOG_INPUT_LINE_MAX_BYTES: usize = 8 * 1024;
|
||||
const AGENT_RUNNER_LOG_OUTPUT_MAX_CHARS: usize = 1_024;
|
||||
const AGENT_RUNNER_CLIENT_EXIT_TIMEOUT: Duration = Duration::from_secs(15);
|
||||
|
||||
#[derive(Default)]
|
||||
pub(super) struct ExternalAgentRunnerGuiOwnerAttachmentState {
|
||||
generation: u64,
|
||||
registration: Option<ExternalAgentRunnerGuiOwnerRegistration>,
|
||||
}
|
||||
|
||||
struct ExternalAgentRunnerGuiOwnerRegistration {
|
||||
generation: u64,
|
||||
config_dir: PathBuf,
|
||||
params: ExternalAgentRunnerRequestParams,
|
||||
attached_boot_id: Option<String>,
|
||||
}
|
||||
|
||||
static EXTERNAL_AGENT_RUNNER_GUI_OWNER_ATTACHMENT_STATE: OnceLock<
|
||||
Mutex<ExternalAgentRunnerGuiOwnerAttachmentState>,
|
||||
> = OnceLock::new();
|
||||
|
||||
fn external_agent_runner_gui_owner_attachment_state(
|
||||
) -> &'static Mutex<ExternalAgentRunnerGuiOwnerAttachmentState> {
|
||||
EXTERNAL_AGENT_RUNNER_GUI_OWNER_ATTACHMENT_STATE
|
||||
.get_or_init(|| Mutex::new(ExternalAgentRunnerGuiOwnerAttachmentState::default()))
|
||||
}
|
||||
|
||||
pub(super) fn register_external_agent_runner_gui_owner_attachment(
|
||||
state: &Mutex<ExternalAgentRunnerGuiOwnerAttachmentState>,
|
||||
config_dir: &Path,
|
||||
params: ExternalAgentRunnerRequestParams,
|
||||
) {
|
||||
let mut state = lock_unpoisoned(state);
|
||||
state.generation = state.generation.wrapping_add(1);
|
||||
let generation = state.generation;
|
||||
state.registration = Some(ExternalAgentRunnerGuiOwnerRegistration {
|
||||
generation,
|
||||
config_dir: config_dir.to_path_buf(),
|
||||
params,
|
||||
attached_boot_id: None,
|
||||
});
|
||||
}
|
||||
|
||||
pub(super) fn attach_registered_external_agent_runner_gui_owner_if_needed_with<F>(
|
||||
state: &Mutex<ExternalAgentRunnerGuiOwnerAttachmentState>,
|
||||
config_dir: &Path,
|
||||
endpoint: &ExternalAgentRunnerEndpoint,
|
||||
attach: F,
|
||||
) -> Result<(), String>
|
||||
where
|
||||
F: FnOnce(&ExternalAgentRunnerEndpoint, ExternalAgentRunnerRequestParams) -> Result<(), String>,
|
||||
{
|
||||
let Some((generation, params)) = ({
|
||||
let state = lock_unpoisoned(state);
|
||||
state.registration.as_ref().and_then(|registration| {
|
||||
(registration.config_dir == config_dir
|
||||
&& registration.attached_boot_id.as_deref() != Some(endpoint.boot_id.as_str()))
|
||||
.then(|| (registration.generation, registration.params.clone()))
|
||||
})
|
||||
}) else {
|
||||
return Ok(());
|
||||
};
|
||||
|
||||
attach(endpoint, params)?;
|
||||
|
||||
let mut state = lock_unpoisoned(state);
|
||||
if let Some(registration) = state.registration.as_mut() {
|
||||
if registration.generation == generation && registration.config_dir == config_dir {
|
||||
registration.attached_boot_id = Some(endpoint.boot_id.clone());
|
||||
}
|
||||
}
|
||||
Ok(())
|
||||
}
|
||||
|
||||
fn redact_url_queries(line: &str) -> String {
|
||||
line.split_whitespace()
|
||||
.map(|token| {
|
||||
@@ -919,18 +990,24 @@ pub(crate) fn attach_external_agent_runner_gui_owner(
|
||||
) -> Result<(), String> {
|
||||
EXTERNAL_AGENT_RUNNER_GUI_OWNER_REQUIRED_CLIENT
|
||||
.store(true, std::sync::atomic::Ordering::Release);
|
||||
let _configure = lock_unpoisoned(external_agent_runner_configure_lock());
|
||||
let config_dir = external_agent_runner_config_dir()
|
||||
.ok_or_else(|| "外部 Agent Runner 尚未配置 AppData;请显式传入 --config-dir".to_string())?;
|
||||
let endpoint = ensure_external_agent_runner(&config_dir)?;
|
||||
let result = send_external_agent_runner_request(
|
||||
&endpoint,
|
||||
"runner.attach_gui_owner",
|
||||
register_external_agent_runner_gui_owner_attachment(
|
||||
external_agent_runner_gui_owner_attachment_state(),
|
||||
&config_dir,
|
||||
ExternalAgentRunnerRequestParams {
|
||||
event_sink_port: Some(event_sink.port),
|
||||
event_sink_token: Some(event_sink.token.clone()),
|
||||
..ExternalAgentRunnerRequestParams::default()
|
||||
},
|
||||
)?;
|
||||
);
|
||||
ensure_external_agent_runner(&config_dir).map(|_| ())
|
||||
}
|
||||
|
||||
pub(super) fn validate_external_agent_runner_gui_owner_attachment_result(
|
||||
result: &Value,
|
||||
) -> Result<(), String> {
|
||||
if result.get("attached").and_then(Value::as_bool) == Some(true)
|
||||
&& result.get("eventSinkAttached").and_then(Value::as_bool) == Some(true)
|
||||
{
|
||||
@@ -940,6 +1017,26 @@ pub(crate) fn attach_external_agent_runner_gui_owner(
|
||||
}
|
||||
}
|
||||
|
||||
fn attach_external_agent_runner_gui_owner_at(
|
||||
endpoint: &ExternalAgentRunnerEndpoint,
|
||||
params: ExternalAgentRunnerRequestParams,
|
||||
) -> Result<(), String> {
|
||||
let result = send_external_agent_runner_request(endpoint, "runner.attach_gui_owner", params)?;
|
||||
validate_external_agent_runner_gui_owner_attachment_result(&result)
|
||||
}
|
||||
|
||||
fn attach_registered_external_agent_runner_gui_owner_if_needed(
|
||||
config_dir: &Path,
|
||||
endpoint: &ExternalAgentRunnerEndpoint,
|
||||
) -> Result<(), String> {
|
||||
attach_registered_external_agent_runner_gui_owner_if_needed_with(
|
||||
external_agent_runner_gui_owner_attachment_state(),
|
||||
config_dir,
|
||||
endpoint,
|
||||
attach_external_agent_runner_gui_owner_at,
|
||||
)
|
||||
}
|
||||
|
||||
pub(super) fn shutdown_external_agent_runner_for_client_exit_at(
|
||||
config_dir: &Path,
|
||||
) -> Result<bool, String> {
|
||||
@@ -1035,6 +1132,9 @@ pub(super) fn ensure_external_agent_runner(
|
||||
match external_agent_runner_endpoint_reuse_decision(&endpoint, &executable_fingerprint) {
|
||||
ExternalAgentRunnerReuseDecision::Reuse => {
|
||||
if ping_external_agent_runner(&endpoint).is_ok() {
|
||||
attach_registered_external_agent_runner_gui_owner_if_needed(
|
||||
config_dir, &endpoint,
|
||||
)?;
|
||||
return Ok(endpoint);
|
||||
}
|
||||
}
|
||||
@@ -1068,6 +1168,7 @@ pub(super) fn ensure_external_agent_runner(
|
||||
let _ = launched.child.wait();
|
||||
})
|
||||
.map_err(|error| format!("启动 Agent Runner 子进程回收线程失败:{error}"))?;
|
||||
attach_registered_external_agent_runner_gui_owner_if_needed(config_dir, &endpoint)?;
|
||||
Ok(endpoint)
|
||||
}
|
||||
Err(error) => {
|
||||
|
||||
@@ -10,6 +10,7 @@ use std::io::{self, Cursor};
|
||||
use std::net::{Ipv4Addr, SocketAddrV4, TcpListener};
|
||||
use std::path::{Path, PathBuf};
|
||||
use std::sync::atomic::{AtomicU64, Ordering};
|
||||
use std::sync::Mutex;
|
||||
use std::time::{Duration, Instant};
|
||||
|
||||
static TEST_DIRECTORY_COUNTER: AtomicU64 = AtomicU64::new(0);
|
||||
@@ -551,6 +552,292 @@ fn runner_endpoint_rejects_hard_links() {
|
||||
assert!(error.contains("硬链接"));
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn gui_owner_registration_replays_once_for_each_runner_boot() {
|
||||
let state = Mutex::new(ExternalAgentRunnerGuiOwnerAttachmentState::default());
|
||||
let config_dir = PathBuf::from("registered-gui-appdata");
|
||||
let event_sink_port = 31_317;
|
||||
let event_sink_token = "a".repeat(64);
|
||||
let params = ExternalAgentRunnerRequestParams {
|
||||
event_sink_port: Some(event_sink_port),
|
||||
event_sink_token: Some(event_sink_token.clone()),
|
||||
..ExternalAgentRunnerRequestParams::default()
|
||||
};
|
||||
register_external_agent_runner_gui_owner_attachment(&state, &config_dir, params);
|
||||
|
||||
let calls = std::cell::RefCell::new(Vec::new());
|
||||
let endpoint_a = test_endpoint(
|
||||
"gui-owner-replay-token-gui-owner-replay-token",
|
||||
"gui-owner-boot-a",
|
||||
31318,
|
||||
);
|
||||
attach_registered_external_agent_runner_gui_owner_if_needed_with(
|
||||
&state,
|
||||
&config_dir,
|
||||
&endpoint_a,
|
||||
|endpoint, params| {
|
||||
calls.borrow_mut().push((
|
||||
endpoint.boot_id.clone(),
|
||||
params
|
||||
.event_sink_port
|
||||
.expect("registered sink port is retained"),
|
||||
params
|
||||
.event_sink_token
|
||||
.expect("registered sink token is retained"),
|
||||
));
|
||||
Ok(())
|
||||
},
|
||||
)
|
||||
.expect("first boot attaches");
|
||||
attach_registered_external_agent_runner_gui_owner_if_needed_with(
|
||||
&state,
|
||||
&config_dir,
|
||||
&endpoint_a,
|
||||
|_, _| panic!("same boot must not attach twice"),
|
||||
)
|
||||
.expect("same boot is idempotent");
|
||||
|
||||
let endpoint_b = test_endpoint(
|
||||
"gui-owner-replay-token-gui-owner-replay-token",
|
||||
"gui-owner-boot-b",
|
||||
31319,
|
||||
);
|
||||
attach_registered_external_agent_runner_gui_owner_if_needed_with(
|
||||
&state,
|
||||
&config_dir,
|
||||
&endpoint_b,
|
||||
|endpoint, params| {
|
||||
calls.borrow_mut().push((
|
||||
endpoint.boot_id.clone(),
|
||||
params
|
||||
.event_sink_port
|
||||
.expect("registered sink port is replayed"),
|
||||
params
|
||||
.event_sink_token
|
||||
.expect("registered sink token is replayed"),
|
||||
));
|
||||
Ok(())
|
||||
},
|
||||
)
|
||||
.expect("replacement boot reattaches");
|
||||
|
||||
assert_eq!(
|
||||
calls.into_inner(),
|
||||
vec![
|
||||
(
|
||||
"gui-owner-boot-a".to_string(),
|
||||
event_sink_port,
|
||||
event_sink_token.clone(),
|
||||
),
|
||||
(
|
||||
"gui-owner-boot-b".to_string(),
|
||||
event_sink_port,
|
||||
event_sink_token,
|
||||
),
|
||||
]
|
||||
);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn gui_owner_registration_failed_replay_remains_pending_for_same_boot() {
|
||||
let state = Mutex::new(ExternalAgentRunnerGuiOwnerAttachmentState::default());
|
||||
let config_dir = PathBuf::from("retry-gui-appdata");
|
||||
register_external_agent_runner_gui_owner_attachment(
|
||||
&state,
|
||||
&config_dir,
|
||||
ExternalAgentRunnerRequestParams::default(),
|
||||
);
|
||||
let endpoint = test_endpoint(
|
||||
"gui-owner-retry-token-gui-owner-retry-token",
|
||||
"gui-owner-retry-boot",
|
||||
31320,
|
||||
);
|
||||
let attempts = std::cell::Cell::new(0_u32);
|
||||
|
||||
let error = attach_registered_external_agent_runner_gui_owner_if_needed_with(
|
||||
&state,
|
||||
&config_dir,
|
||||
&endpoint,
|
||||
|_, _| {
|
||||
attempts.set(attempts.get() + 1);
|
||||
Err("injected attach failure".to_string())
|
||||
},
|
||||
)
|
||||
.expect_err("failed attach must remain pending");
|
||||
assert_eq!(error, "injected attach failure");
|
||||
attach_registered_external_agent_runner_gui_owner_if_needed_with(
|
||||
&state,
|
||||
&config_dir,
|
||||
&endpoint,
|
||||
|_, _| {
|
||||
attempts.set(attempts.get() + 1);
|
||||
Ok(())
|
||||
},
|
||||
)
|
||||
.expect("same boot retries after failure");
|
||||
attach_registered_external_agent_runner_gui_owner_if_needed_with(
|
||||
&state,
|
||||
&config_dir,
|
||||
&endpoint,
|
||||
|_, _| panic!("successful retry must mark the boot attached"),
|
||||
)
|
||||
.expect("successful retry is idempotent");
|
||||
assert_eq!(attempts.get(), 2);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn gui_owner_registration_missing_event_sink_confirmation_retries_same_boot() {
|
||||
let state = Mutex::new(ExternalAgentRunnerGuiOwnerAttachmentState::default());
|
||||
let config_dir = PathBuf::from("missing-sink-confirmation-appdata");
|
||||
register_external_agent_runner_gui_owner_attachment(
|
||||
&state,
|
||||
&config_dir,
|
||||
ExternalAgentRunnerRequestParams {
|
||||
event_sink_port: Some(31_322),
|
||||
event_sink_token: Some("c".repeat(64)),
|
||||
..ExternalAgentRunnerRequestParams::default()
|
||||
},
|
||||
);
|
||||
let endpoint = test_endpoint(
|
||||
"missing-sink-confirmation-runner-token",
|
||||
"missing-sink-confirmation-boot",
|
||||
31_322,
|
||||
);
|
||||
let attempts = std::cell::Cell::new(0_u32);
|
||||
|
||||
attach_registered_external_agent_runner_gui_owner_if_needed_with(
|
||||
&state,
|
||||
&config_dir,
|
||||
&endpoint,
|
||||
|_, _| {
|
||||
attempts.set(attempts.get() + 1);
|
||||
validate_external_agent_runner_gui_owner_attachment_result(&json!({
|
||||
"attached": true
|
||||
}))
|
||||
},
|
||||
)
|
||||
.expect_err("missing eventSinkAttached must fail");
|
||||
attach_registered_external_agent_runner_gui_owner_if_needed_with(
|
||||
&state,
|
||||
&config_dir,
|
||||
&endpoint,
|
||||
|_, _| {
|
||||
attempts.set(attempts.get() + 1);
|
||||
validate_external_agent_runner_gui_owner_attachment_result(&json!({
|
||||
"attached": true,
|
||||
"eventSinkAttached": true
|
||||
}))
|
||||
},
|
||||
)
|
||||
.expect("same boot retries after missing event sink confirmation");
|
||||
assert_eq!(attempts.get(), 2);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn gui_owner_registration_false_event_sink_confirmation_retries_same_boot() {
|
||||
let state = Mutex::new(ExternalAgentRunnerGuiOwnerAttachmentState::default());
|
||||
let config_dir = PathBuf::from("false-sink-confirmation-appdata");
|
||||
register_external_agent_runner_gui_owner_attachment(
|
||||
&state,
|
||||
&config_dir,
|
||||
ExternalAgentRunnerRequestParams {
|
||||
event_sink_port: Some(31_323),
|
||||
event_sink_token: Some("d".repeat(64)),
|
||||
..ExternalAgentRunnerRequestParams::default()
|
||||
},
|
||||
);
|
||||
let endpoint = test_endpoint(
|
||||
"false-sink-confirmation-runner-token",
|
||||
"false-sink-confirmation-boot",
|
||||
31_323,
|
||||
);
|
||||
let attempts = std::cell::Cell::new(0_u32);
|
||||
|
||||
attach_registered_external_agent_runner_gui_owner_if_needed_with(
|
||||
&state,
|
||||
&config_dir,
|
||||
&endpoint,
|
||||
|_, _| {
|
||||
attempts.set(attempts.get() + 1);
|
||||
validate_external_agent_runner_gui_owner_attachment_result(&json!({
|
||||
"attached": true,
|
||||
"eventSinkAttached": false
|
||||
}))
|
||||
},
|
||||
)
|
||||
.expect_err("false eventSinkAttached must fail");
|
||||
attach_registered_external_agent_runner_gui_owner_if_needed_with(
|
||||
&state,
|
||||
&config_dir,
|
||||
&endpoint,
|
||||
|_, _| {
|
||||
attempts.set(attempts.get() + 1);
|
||||
validate_external_agent_runner_gui_owner_attachment_result(&json!({
|
||||
"attached": true,
|
||||
"eventSinkAttached": true
|
||||
}))
|
||||
},
|
||||
)
|
||||
.expect("same boot retries after false event sink confirmation");
|
||||
assert_eq!(attempts.get(), 2);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn gui_owner_registration_does_not_cross_config_dirs() {
|
||||
let state = Mutex::new(ExternalAgentRunnerGuiOwnerAttachmentState::default());
|
||||
let registered_config_dir = PathBuf::from("registered-gui-appdata");
|
||||
let other_config_dir = PathBuf::from("other-gui-appdata");
|
||||
let event_sink_token = "e".repeat(64);
|
||||
register_external_agent_runner_gui_owner_attachment(
|
||||
&state,
|
||||
®istered_config_dir,
|
||||
ExternalAgentRunnerRequestParams {
|
||||
event_sink_port: Some(31_324),
|
||||
event_sink_token: Some(event_sink_token.clone()),
|
||||
..ExternalAgentRunnerRequestParams::default()
|
||||
},
|
||||
);
|
||||
let endpoint = test_endpoint(
|
||||
"gui-owner-config-token-gui-owner-config-token",
|
||||
"gui-owner-config-boot",
|
||||
31321,
|
||||
);
|
||||
attach_registered_external_agent_runner_gui_owner_if_needed_with(
|
||||
&state,
|
||||
&other_config_dir,
|
||||
&endpoint,
|
||||
|_, _| panic!("GUI owner registration must stay bound to its AppData"),
|
||||
)
|
||||
.expect("other AppData remains unattached");
|
||||
|
||||
let calls = std::cell::Cell::new(0_u32);
|
||||
attach_registered_external_agent_runner_gui_owner_if_needed_with(
|
||||
&state,
|
||||
®istered_config_dir,
|
||||
&endpoint,
|
||||
|_, params| {
|
||||
calls.set(calls.get() + 1);
|
||||
assert_eq!(params.event_sink_port, Some(31_324));
|
||||
assert_eq!(
|
||||
params.event_sink_token.as_deref(),
|
||||
Some(event_sink_token.as_str())
|
||||
);
|
||||
Ok(())
|
||||
},
|
||||
)
|
||||
.expect("registered AppData attaches");
|
||||
assert_eq!(calls.get(), 1);
|
||||
|
||||
let unregistered = Mutex::new(ExternalAgentRunnerGuiOwnerAttachmentState::default());
|
||||
attach_registered_external_agent_runner_gui_owner_if_needed_with(
|
||||
&unregistered,
|
||||
®istered_config_dir,
|
||||
&endpoint,
|
||||
|_, _| panic!("CLI state without GUI registration must not attach"),
|
||||
)
|
||||
.expect("unregistered CLI state remains unchanged");
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn gui_owner_lock_allows_only_one_frontend_process_per_appdata() {
|
||||
let directory = unique_test_directory();
|
||||
@@ -567,7 +854,9 @@ fn gui_owner_lock_allows_only_one_frontend_process_per_appdata() {
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn attached_gui_owner_loss_forces_runner_shutdown() {
|
||||
fn manifest_invalidation_sink_isolation_gui_owner_attach_configures_and_cleans_up() {
|
||||
let sink_guard = crate::acquire_game_creator_manifest_invalidation_event_sink_test_guard();
|
||||
assert_eq!(sink_guard.configured_sink(), None);
|
||||
let directory = unique_test_directory();
|
||||
let config_dir = private_runner_test_config_dir(&directory);
|
||||
let token = "gui-owner-monitor-token-gui-owner-monitor-token";
|
||||
@@ -593,6 +882,13 @@ fn attached_gui_owner_loss_forces_runner_shutdown() {
|
||||
);
|
||||
assert!(attached.ok);
|
||||
assert!(state.gui_owner_attached.load(Ordering::Acquire));
|
||||
assert_eq!(
|
||||
sink_guard.configured_sink(),
|
||||
Some(crate::GameCreatorManifestInvalidationEventSink {
|
||||
port: 31_318,
|
||||
token: "b".repeat(64),
|
||||
})
|
||||
);
|
||||
assert!(
|
||||
!external_agent_runner_shutdown_if_gui_owner_lost(&state).expect("owner remains present")
|
||||
);
|
||||
@@ -603,7 +899,9 @@ fn attached_gui_owner_loss_forces_runner_shutdown() {
|
||||
assert!(state.draining.load(Ordering::Acquire));
|
||||
assert!(state.force_shutdown_requested.load(Ordering::Acquire));
|
||||
assert!(state.shutdown_requested.load(Ordering::Acquire));
|
||||
crate::clear_game_creator_manifest_invalidation_event_sink_for_test();
|
||||
drop(sink_guard);
|
||||
let cleanup_guard = crate::acquire_game_creator_manifest_invalidation_event_sink_test_guard();
|
||||
assert_eq!(cleanup_guard.configured_sink(), None);
|
||||
}
|
||||
|
||||
#[test]
|
||||
|
||||
@@ -1606,11 +1606,9 @@ async fn project_supervisor_resume_rechecks_delegate_policy_after_delivery_reser
|
||||
.expect("read barrier after rejecting reserved delivery")
|
||||
.is_clear());
|
||||
|
||||
let released = wait_for_agent_runtime_lane_release_async(
|
||||
&root,
|
||||
GAME_CREATOR_PROJECT_SUPERVISOR_AGENT_ID,
|
||||
)
|
||||
.await;
|
||||
let released =
|
||||
wait_for_agent_runtime_lane_release_async(&root, GAME_CREATOR_PROJECT_SUPERVISOR_AGENT_ID)
|
||||
.await;
|
||||
assert_eq!(released.state.run_id, parent_run_id);
|
||||
|
||||
fs::remove_dir_all(root).ok();
|
||||
|
||||
@@ -3,18 +3,75 @@ use base64::Engine as _;
|
||||
use serde_json::Value;
|
||||
use sha2::{Digest as _, Sha256};
|
||||
use std::collections::{BTreeMap, BTreeSet};
|
||||
use std::io::{Read, Write};
|
||||
use std::io::{self, Read, Write};
|
||||
use std::sync::atomic::{AtomicBool, AtomicU64, Ordering};
|
||||
use std::sync::{Arc, Barrier, Condvar, Mutex as StdMutex, MutexGuard as StdMutexGuard};
|
||||
use std::time::{SystemTime, UNIX_EPOCH};
|
||||
use std::time::{Instant, SystemTime, UNIX_EPOCH};
|
||||
use zip::write::SimpleFileOptions;
|
||||
|
||||
static TEST_PROJECT_COUNTER: AtomicU64 = AtomicU64::new(0);
|
||||
static TEST_MOCK_PORT_COUNTER: AtomicU64 = AtomicU64::new(20_000);
|
||||
static TEST_CONFIG_LOCK: StdMutex<()> = StdMutex::new(());
|
||||
const MANIFEST_INVALIDATION_RELAY_TEST_ACCEPT_TIMEOUT: Duration = Duration::from_millis(500);
|
||||
const MANIFEST_INVALIDATION_RELAY_TEST_PAYLOAD_TIMEOUT: Duration = Duration::from_millis(500);
|
||||
const MANIFEST_INVALIDATION_RELAY_TEST_MAX_BYTES: usize = 64 * 1024;
|
||||
|
||||
fn read_manifest_invalidation_relay_payload_with_deadline(
|
||||
listener: &TcpListener,
|
||||
) -> io::Result<Vec<u8>> {
|
||||
listener.set_nonblocking(true)?;
|
||||
let accept_deadline = Instant::now() + MANIFEST_INVALIDATION_RELAY_TEST_ACCEPT_TIMEOUT;
|
||||
let (mut stream, _) = loop {
|
||||
match listener.accept() {
|
||||
Ok(accepted) => break accepted,
|
||||
Err(error) if error.kind() == io::ErrorKind::WouldBlock => {
|
||||
if Instant::now() >= accept_deadline {
|
||||
return Err(io::Error::new(
|
||||
io::ErrorKind::TimedOut,
|
||||
"manifest invalidation relay accept timed out",
|
||||
));
|
||||
}
|
||||
std::thread::yield_now();
|
||||
}
|
||||
Err(error) if error.kind() == io::ErrorKind::Interrupted => continue,
|
||||
Err(error) => return Err(error),
|
||||
}
|
||||
};
|
||||
|
||||
stream.set_nonblocking(true)?;
|
||||
let payload_deadline = Instant::now() + MANIFEST_INVALIDATION_RELAY_TEST_PAYLOAD_TIMEOUT;
|
||||
let mut payload = Vec::new();
|
||||
let mut buffer = [0_u8; 4096];
|
||||
loop {
|
||||
match stream.read(&mut buffer) {
|
||||
Ok(0) => return Ok(payload),
|
||||
Ok(read) => {
|
||||
payload.extend_from_slice(&buffer[..read]);
|
||||
if payload.len() > MANIFEST_INVALIDATION_RELAY_TEST_MAX_BYTES {
|
||||
return Err(io::Error::new(
|
||||
io::ErrorKind::InvalidData,
|
||||
"manifest invalidation relay payload exceeded test limit",
|
||||
));
|
||||
}
|
||||
}
|
||||
Err(error) if error.kind() == io::ErrorKind::WouldBlock => {
|
||||
if Instant::now() >= payload_deadline {
|
||||
return Err(io::Error::new(
|
||||
io::ErrorKind::TimedOut,
|
||||
"manifest invalidation relay payload timed out",
|
||||
));
|
||||
}
|
||||
std::thread::yield_now();
|
||||
}
|
||||
Err(error) if error.kind() == io::ErrorKind::Interrupted => continue,
|
||||
Err(error) => return Err(error),
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn non_supervisor_runtime_update_invalidates_manifest_on_the_wire_and_runner_relay() {
|
||||
fn manifest_invalidation_sink_isolation_relays_non_supervisor_runtime_update() {
|
||||
let sink_guard = acquire_game_creator_manifest_invalidation_event_sink_test_guard();
|
||||
let root = unique_project_path();
|
||||
init_local_game_project_at(&root, "runtime-event-contract", "Runtime 事件合同测试")
|
||||
.expect("init runtime event contract project");
|
||||
@@ -33,20 +90,12 @@ fn non_supervisor_runtime_update_invalidates_manifest_on_the_wire_and_runner_rel
|
||||
.expect("read manifest invalidation relay fixture address")
|
||||
.port();
|
||||
let relay_token = "a".repeat(64);
|
||||
configure_game_creator_manifest_invalidation_event_sink(relay_port, &relay_token)
|
||||
sink_guard
|
||||
.configure(relay_port, &relay_token)
|
||||
.expect("configure manifest invalidation relay fixture");
|
||||
emit_game_creator_agent_runtime_update(&root, "art-asset-plan");
|
||||
let (mut relay_stream, _) = relay_listener
|
||||
.accept()
|
||||
.expect("accept manifest invalidation relay");
|
||||
relay_stream
|
||||
.set_read_timeout(Some(Duration::from_secs(1)))
|
||||
.expect("set manifest invalidation relay read timeout");
|
||||
let mut relay_payload = Vec::new();
|
||||
relay_stream
|
||||
.read_to_end(&mut relay_payload)
|
||||
.expect("read manifest invalidation relay");
|
||||
clear_game_creator_manifest_invalidation_event_sink_for_test();
|
||||
let relay_payload = read_manifest_invalidation_relay_payload_with_deadline(&relay_listener)
|
||||
.expect("receive manifest invalidation relay within deadline");
|
||||
let relay: GameCreatorManifestInvalidationRelayEnvelope =
|
||||
serde_json::from_slice(&relay_payload).expect("parse manifest invalidation relay");
|
||||
assert_eq!(relay.token, relay_token);
|
||||
@@ -56,6 +105,58 @@ fn non_supervisor_runtime_update_invalidates_manifest_on_the_wire_and_runner_rel
|
||||
fs::remove_dir_all(root).ok();
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn manifest_invalidation_sink_isolation_bounds_timeouts_and_cleans_up_with_raii() {
|
||||
let cleanup_listener = TcpListener::bind((std::net::Ipv4Addr::LOCALHOST, 0))
|
||||
.expect("bind manifest invalidation cleanup fixture");
|
||||
let cleanup_port = cleanup_listener
|
||||
.local_addr()
|
||||
.expect("read manifest invalidation cleanup fixture address")
|
||||
.port();
|
||||
let cleanup_token = "b".repeat(64);
|
||||
let unwind = std::panic::catch_unwind(|| {
|
||||
let sink_guard = acquire_game_creator_manifest_invalidation_event_sink_test_guard();
|
||||
sink_guard
|
||||
.configure(cleanup_port, &cleanup_token)
|
||||
.expect("configure manifest invalidation cleanup fixture");
|
||||
assert_eq!(
|
||||
sink_guard.configured_sink(),
|
||||
Some(GameCreatorManifestInvalidationEventSink {
|
||||
port: cleanup_port,
|
||||
token: cleanup_token.clone(),
|
||||
})
|
||||
);
|
||||
panic!("exercise manifest invalidation sink guard unwind cleanup");
|
||||
});
|
||||
assert!(unwind.is_err());
|
||||
|
||||
let sink_guard = acquire_game_creator_manifest_invalidation_event_sink_test_guard();
|
||||
assert_eq!(sink_guard.configured_sink(), None);
|
||||
|
||||
let empty_listener = TcpListener::bind((std::net::Ipv4Addr::LOCALHOST, 0))
|
||||
.expect("bind empty manifest invalidation relay fixture");
|
||||
let accept_started = Instant::now();
|
||||
let accept_error = read_manifest_invalidation_relay_payload_with_deadline(&empty_listener)
|
||||
.expect_err("missing relay must time out");
|
||||
assert_eq!(accept_error.kind(), io::ErrorKind::TimedOut);
|
||||
assert!(accept_started.elapsed() < Duration::from_secs(2));
|
||||
|
||||
let stalled_listener = TcpListener::bind((std::net::Ipv4Addr::LOCALHOST, 0))
|
||||
.expect("bind stalled manifest invalidation relay fixture");
|
||||
let stalled_stream = TcpStream::connect(
|
||||
stalled_listener
|
||||
.local_addr()
|
||||
.expect("read stalled manifest invalidation relay fixture address"),
|
||||
)
|
||||
.expect("connect stalled manifest invalidation relay fixture");
|
||||
let payload_started = Instant::now();
|
||||
let payload_error = read_manifest_invalidation_relay_payload_with_deadline(&stalled_listener)
|
||||
.expect_err("incomplete relay payload must time out");
|
||||
assert_eq!(payload_error.kind(), io::ErrorKind::TimedOut);
|
||||
assert!(payload_started.elapsed() < Duration::from_secs(2));
|
||||
drop(stalled_stream);
|
||||
}
|
||||
|
||||
#[test]
|
||||
fn gui_final_exit_is_the_only_run_event_that_requests_runner_shutdown() {
|
||||
assert!(game_creator_gui_run_event_requests_runner_shutdown(
|
||||
|
||||
@@ -3,6 +3,7 @@
|
||||
import './assetCanvasSurface.css';
|
||||
|
||||
import {
|
||||
type CanvasHistoryAction,
|
||||
type CanvasLayer,
|
||||
type CanvasViewport,
|
||||
createMinimapModel,
|
||||
@@ -13,7 +14,6 @@ import {
|
||||
type ImageCanvasGenerationProgressPhase,
|
||||
type ImageCanvasHostScope,
|
||||
type ImageCanvasMediaRef,
|
||||
MAX_HISTORY_STEPS,
|
||||
moveViewportFromMinimapPointer,
|
||||
moveViewportFromPan,
|
||||
removeCanvasLayers,
|
||||
@@ -26,6 +26,7 @@ import {
|
||||
CanvasWorld,
|
||||
LayerRenderer,
|
||||
Minimap,
|
||||
useCanvasHistory,
|
||||
ZoomControls,
|
||||
} from '@genarrative/image-canvas-react';
|
||||
import {
|
||||
@@ -85,11 +86,6 @@ export type AssetCanvasCommitNotification = {
|
||||
};
|
||||
|
||||
type RuntimeCanvasLayer = CanvasLayer & { mediaRef: ImageCanvasMediaRef };
|
||||
type HistorySnapshot = {
|
||||
layers: RuntimeCanvasLayer[];
|
||||
viewport: CanvasViewport;
|
||||
selectedLayerIds: string[];
|
||||
};
|
||||
|
||||
type PendingGenerationIdentity = {
|
||||
saveAttemptId: string;
|
||||
@@ -232,14 +228,6 @@ function draftCanvasFromRuntime(
|
||||
};
|
||||
}
|
||||
|
||||
function historyClone(input: HistorySnapshot): HistorySnapshot {
|
||||
return {
|
||||
layers: input.layers.map((layer) => ({ ...layer })),
|
||||
viewport: { ...input.viewport },
|
||||
selectedLayerIds: [...input.selectedLayerIds],
|
||||
};
|
||||
}
|
||||
|
||||
function mediaTypeForFile(file: File) {
|
||||
if (file.type === 'image/png') return 'image/png' as const;
|
||||
if (file.type === 'image/jpeg') return 'image/jpeg' as const;
|
||||
@@ -317,8 +305,6 @@ export function AssetCanvasSurface({
|
||||
const documentVersionRef = useRef(documentVersion);
|
||||
const epochRef = useRef(0);
|
||||
const dragRef = useRef<DragState | null>(null);
|
||||
const undoRef = useRef<HistorySnapshot[]>([]);
|
||||
const redoRef = useRef<HistorySnapshot[]>([]);
|
||||
const saveQueueRef = useRef<Promise<unknown>>(Promise.resolve());
|
||||
const savePromiseRef = useRef<Promise<void> | null>(null);
|
||||
const hostRevisionRef = useRef(expectedHostRevision);
|
||||
@@ -346,54 +332,45 @@ export function AssetCanvasSurface({
|
||||
setLifecycle({ kind: 'canvas.editing', dirty: true });
|
||||
}, []);
|
||||
|
||||
const currentSnapshot = useCallback(
|
||||
(): HistorySnapshot => ({
|
||||
layers: layersRef.current.map((layer) => ({ ...layer })),
|
||||
viewport: { ...viewportRef.current },
|
||||
selectedLayerIds: [...selectionRef.current],
|
||||
const canvasHistoryRefs = useMemo(
|
||||
() => ({
|
||||
layersRef,
|
||||
viewportRef,
|
||||
selectedLayerIdsRef: selectionRef,
|
||||
}),
|
||||
[],
|
||||
);
|
||||
|
||||
const captureHistory = useCallback(() => {
|
||||
undoRef.current = [
|
||||
...undoRef.current.slice(-(MAX_HISTORY_STEPS - 1)),
|
||||
historyClone(currentSnapshot()),
|
||||
];
|
||||
redoRef.current = [];
|
||||
}, [currentSnapshot]);
|
||||
|
||||
const applySnapshot = useCallback(
|
||||
(snapshot: HistorySnapshot) => {
|
||||
setLayers(snapshot.layers.map((layer) => ({ ...layer })));
|
||||
setViewport({ ...snapshot.viewport });
|
||||
setSelectedLayerIds([...snapshot.selectedLayerIds]);
|
||||
markDirty();
|
||||
},
|
||||
[markDirty],
|
||||
const canvasHistorySetters = useMemo(
|
||||
() => ({
|
||||
setLayers: (nextLayers: CanvasLayer[]) =>
|
||||
setLayers(nextLayers as RuntimeCanvasLayer[]),
|
||||
setViewport,
|
||||
setSelectedLayerIds,
|
||||
}),
|
||||
[],
|
||||
);
|
||||
const {
|
||||
canUndo,
|
||||
canRedo,
|
||||
captureCanvasHistory,
|
||||
undoCanvasChange,
|
||||
redoCanvasChange,
|
||||
resetCanvasHistory,
|
||||
} = useCanvasHistory({
|
||||
refs: canvasHistoryRefs,
|
||||
setters: canvasHistorySetters,
|
||||
allowContentRemovalOnRestore: true,
|
||||
});
|
||||
const captureHistory = useCallback(
|
||||
(action: CanvasHistoryAction) => captureCanvasHistory(action),
|
||||
[captureCanvasHistory],
|
||||
);
|
||||
|
||||
const undo = useCallback(() => {
|
||||
const previous = undoRef.current.at(-1);
|
||||
if (!previous) return;
|
||||
redoRef.current = [
|
||||
...redoRef.current.slice(-(MAX_HISTORY_STEPS - 1)),
|
||||
historyClone(currentSnapshot()),
|
||||
];
|
||||
undoRef.current = undoRef.current.slice(0, -1);
|
||||
applySnapshot(previous);
|
||||
}, [applySnapshot, currentSnapshot]);
|
||||
|
||||
if (undoCanvasChange().status === 'success') markDirty();
|
||||
}, [markDirty, undoCanvasChange]);
|
||||
const redo = useCallback(() => {
|
||||
const next = redoRef.current.at(-1);
|
||||
if (!next) return;
|
||||
undoRef.current = [
|
||||
...undoRef.current.slice(-(MAX_HISTORY_STEPS - 1)),
|
||||
historyClone(currentSnapshot()),
|
||||
];
|
||||
redoRef.current = redoRef.current.slice(0, -1);
|
||||
applySnapshot(next);
|
||||
}, [applySnapshot, currentSnapshot]);
|
||||
if (redoCanvasChange().status === 'success') markDirty();
|
||||
}, [markDirty, redoCanvasChange]);
|
||||
|
||||
const hydrateDraft = useCallback(
|
||||
async (nextDraft: ImageCanvasDraft, epoch: number) => {
|
||||
@@ -446,11 +423,10 @@ export function AssetCanvasSurface({
|
||||
setViewport(nextDraft.canvas.viewport);
|
||||
setBackgroundColor(nextDraft.canvas.backgroundColor);
|
||||
setSelectedLayerIds(nextDraft.canvas.selectedLayerIds);
|
||||
undoRef.current = [];
|
||||
redoRef.current = [];
|
||||
resetCanvasHistory();
|
||||
setLifecycle({ kind: 'canvas.editing', dirty: false });
|
||||
},
|
||||
[host, stableScope],
|
||||
[host, resetCanvasHistory, stableScope],
|
||||
);
|
||||
|
||||
useEffect(() => {
|
||||
@@ -784,7 +760,7 @@ export function AssetCanvasSurface({
|
||||
);
|
||||
return;
|
||||
}
|
||||
captureHistory();
|
||||
captureHistory({ type: 'upload-image', count: files.length });
|
||||
const baseZ = layersRef.current.reduce(
|
||||
(value, layer) => Math.max(value, layer.zIndex),
|
||||
-1,
|
||||
@@ -838,7 +814,11 @@ export function AssetCanvasSurface({
|
||||
|
||||
const deleteSelected = useCallback(() => {
|
||||
if (!selectionRef.current.length) return;
|
||||
captureHistory();
|
||||
captureHistory({
|
||||
type: 'delete-image',
|
||||
count: selectionRef.current.length,
|
||||
layerIds: [...selectionRef.current],
|
||||
});
|
||||
setLayers(
|
||||
(current) =>
|
||||
removeCanvasLayers(
|
||||
@@ -1304,18 +1284,14 @@ export function AssetCanvasSurface({
|
||||
<button
|
||||
type="button"
|
||||
onClick={undo}
|
||||
disabled={
|
||||
!undoRef.current.length || lifecycle.kind === 'canvas.generating'
|
||||
}
|
||||
disabled={!canUndo || lifecycle.kind === 'canvas.generating'}
|
||||
>
|
||||
撤销
|
||||
</button>
|
||||
<button
|
||||
type="button"
|
||||
onClick={redo}
|
||||
disabled={
|
||||
!redoRef.current.length || lifecycle.kind === 'canvas.generating'
|
||||
}
|
||||
disabled={!canRedo || lifecycle.kind === 'canvas.generating'}
|
||||
>
|
||||
重做
|
||||
</button>
|
||||
@@ -1527,7 +1503,7 @@ export function AssetCanvasSurface({
|
||||
isPanning={dragRef.current?.kind === 'pan'}
|
||||
onPointerDown={(event) => {
|
||||
if (event.target !== event.currentTarget) return;
|
||||
captureHistory();
|
||||
captureHistory({ type: 'change-viewport' });
|
||||
setSelectedLayerIds([]);
|
||||
dragRef.current = {
|
||||
kind: 'pan',
|
||||
@@ -1583,7 +1559,11 @@ export function AssetCanvasSurface({
|
||||
}
|
||||
return;
|
||||
}
|
||||
captureHistory();
|
||||
captureHistory({
|
||||
type: 'move-image',
|
||||
count: targetIds.length,
|
||||
layerIds: targetIds,
|
||||
});
|
||||
dragRef.current = {
|
||||
kind: 'move',
|
||||
startClientX: event.clientX,
|
||||
@@ -1613,7 +1593,11 @@ export function AssetCanvasSurface({
|
||||
className="asset-canvas-surface__resize-handle"
|
||||
onPointerDown={(event) => {
|
||||
event.stopPropagation();
|
||||
captureHistory();
|
||||
captureHistory({
|
||||
type: 'resize-image',
|
||||
count: 1,
|
||||
layerIds: [layer.id],
|
||||
});
|
||||
dragRef.current = {
|
||||
kind: 'resize',
|
||||
startClientX: event.clientX,
|
||||
@@ -1638,13 +1622,13 @@ export function AssetCanvasSurface({
|
||||
onFit={() => {
|
||||
const next = fitViewportToLayers({ layers, canvasSize });
|
||||
if (next) {
|
||||
captureHistory();
|
||||
captureHistory({ type: 'change-viewport' });
|
||||
setViewport(next);
|
||||
markDirty();
|
||||
}
|
||||
}}
|
||||
onScaleFromCenter={(scale) => {
|
||||
captureHistory();
|
||||
captureHistory({ type: 'change-viewport' });
|
||||
setViewport((current) =>
|
||||
scaleViewportFromScreenPoint({
|
||||
viewport: current,
|
||||
@@ -1685,7 +1669,7 @@ export function AssetCanvasSurface({
|
||||
model={minimapModel}
|
||||
onPointerDown={(event) => {
|
||||
const rect = event.currentTarget.getBoundingClientRect();
|
||||
captureHistory();
|
||||
captureHistory({ type: 'change-viewport' });
|
||||
setViewport(
|
||||
moveViewportFromMinimapPointer({
|
||||
viewport,
|
||||
|
||||
@@ -1,9 +1,9 @@
|
||||
import type {
|
||||
ImageCanvasDraft,
|
||||
ImageCanvasGenerationCommitResult,
|
||||
ImageCanvasGenerationPort,
|
||||
ImageCanvasGenerationProgress,
|
||||
ImageCanvasGenerationRecord,
|
||||
ImageCanvasGenerationPort,
|
||||
ImageCanvasHostPort,
|
||||
ImageCanvasHostResult,
|
||||
ImageCanvasHostScope,
|
||||
|
||||
Reference in New Issue
Block a user