Files
Genarrative/apps/ai-game-creator-shell/src-tauri/src/agent/direct_execution.rs
T
k88936 cd196cf61c 宿主:开发构建跳过 Codex 执行器版本门禁
- codex_app_server 逐次审批门禁与 direct_execution 补丁执行器门禁改为按 profile 分流:发行构建仍要求严格等于捆绑侧车固定版本,开发构建(debug_assertions)直接通过
- 修正开发态必然被拒的问题:开发构建从宿主 PATH 解析到的 Codex(本机 codex-cli 0.156.0)与固定版本 codex-cli 0.155.1 不等,且 Linux 与未 stage 侧车时没有可选固定版本,导致 Direct 回合在建连前就被拒
- 发行构建的拒单文案补上期望版本与实际版本,便于排障
- 同步调整受影响的单测:开发构建断言跳过门禁,发行构建断言仍拒绝版本漂移
2026-09-24 17:42:25 +08:00

1300 lines
48 KiB
Rust
Raw Blame History

This file contains ambiguous Unicode characters
This file contains Unicode characters that might be confused with other characters. If you think that this is intentional, you can safely ignore this warning. Use the Escape button to reveal them.
//! Direct 回合的宿主权威执行状态。持久化先于放行,磁盘路径不进入项目或模型上下文。
use super::direct_validation::DirectValidationConfig;
use serde::{Deserialize, Serialize};
use serde_json::{json, Value};
use sha2::{Digest, Sha256};
use std::collections::BTreeMap;
use std::fs::{File, OpenOptions};
use std::path::{Path, PathBuf};
use std::sync::{Arc, Mutex, OnceLock, Weak};
use std::time::{Instant, SystemTime, UNIX_EPOCH};
const SCHEMA: &str = "agc-direct-execution.v1";
const MAX_STATE_BYTES: usize = 1024 * 1024;
const MAX_ACTIVE: usize = 64;
const MAX_EVIDENCE: usize = 64;
static SESSIONS: OnceLock<Mutex<BTreeMap<PathBuf, Weak<ExecutionSession>>>> = OnceLock::new();
#[derive(Clone, Copy, Debug, Deserialize, Serialize, PartialEq, Eq)]
#[serde(rename_all = "kebab-case")]
pub(super) enum ExecutionPhase {
Working,
Draining,
Sealing,
Completed,
Exhausted,
Interrupted,
}
impl ExecutionPhase {
pub(super) fn is_terminal(self) -> bool {
matches!(self, Self::Completed | Self::Exhausted | Self::Interrupted)
}
}
#[derive(Clone, Copy, Debug, Deserialize, Serialize, PartialEq, Eq)]
#[serde(rename_all = "kebab-case")]
pub(super) enum EffectKind {
Write,
Execute,
Paid,
}
#[derive(Clone, Debug, Deserialize, Serialize)]
#[serde(rename_all = "camelCase", deny_unknown_fields)]
pub(super) struct LeaseRecord {
pub(super) kind: EffectKind,
pub(super) sequence: u32,
pub(super) pass: u32,
started_at_ms: u64,
pub(super) validation_source: Option<String>,
pub(super) validation_key: Option<String>,
#[serde(default)]
admitted_revision: u64,
#[serde(default)]
paid_dispatch_count: u32,
}
#[derive(Clone, Debug, Deserialize, Serialize)]
#[serde(rename_all = "camelCase", deny_unknown_fields)]
pub(super) struct ExecutionEvidence {
pub(super) key: String,
pub(super) fingerprint: String,
pub(super) source_fingerprint: String,
pub(super) result: Value,
}
#[derive(Clone, Debug, Deserialize, Serialize)]
#[serde(rename_all = "camelCase", deny_unknown_fields)]
pub(super) struct ExecutionLedger {
schema_version: String,
pub(super) client_turn_id: String,
pub(super) project_id: String,
project_key: String,
request_hash: String,
pub(super) requires_contract: bool,
pub(super) contract: Option<Value>,
pub(super) phase: ExecutionPhase,
pub(super) revision: u64,
pub(super) used_passes: u32,
pub(super) used_execution_ms: u64,
pub(super) max_runs: u32,
pub(super) max_execution_ms: u64,
pub(super) max_turn_ms: u64,
pub(super) created_at_ms: u64,
next_sequence: u32,
validation_source: Option<String>,
pub(super) active: BTreeMap<String, LeaseRecord>,
pub(super) evidence: BTreeMap<String, ExecutionEvidence>,
pub(super) executor_stopped: bool,
pub(super) terminal_report: Option<String>,
#[serde(default)]
pub(super) delivery_reviews: u32,
#[serde(default)]
pub(super) plan: Option<Value>,
#[serde(default)]
pub(super) last_failed_write_revision: Option<u64>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub(super) analytics_run: Option<crate::analytics::run::Metadata>,
}
struct SessionData {
ledger: ExecutionLedger,
started: Instant,
initial_elapsed_ms: u64,
lease_started: BTreeMap<String, Instant>,
poisoned: bool,
codex_executor: Option<CodexExecutorIdentity>,
#[cfg(test)]
elapsed_offset_ms: u64,
}
// 与业务持久化锁分离;只保存内存状态,持锁期间不执行 I/O 或投递事件。
struct SessionAnalytics {
project_id: String,
route: Option<crate::analytics::contract::Route>,
capture: Option<(
crate::analytics::contract::Context,
crate::analytics::store::AnalyticsWriter,
)>,
output_revision: Option<u64>,
}
#[derive(Clone, PartialEq, Eq)]
struct CodexExecutorIdentity {
path: PathBuf,
version: String,
digest: String,
}
fn executor_digest(path: &Path) -> Result<String, String> {
use std::io::Read;
let mut file =
File::open(path).map_err(|_| "direct-execution-executor: 无法读取已选择的执行器")?;
let before = file
.metadata()
.map_err(|_| "direct-execution-executor: 执行器身份不可用")?;
if !before.is_file() || before.len() > 512 * 1024 * 1024 {
return Err("direct-execution-executor: 执行器类型或大小无效".into());
}
let mut digest = Sha256::new();
let mut bytes = 0u64;
let mut buffer = [0u8; 65536];
loop {
let read = file
.read(&mut buffer)
.map_err(|_| "direct-execution-executor: 执行器读取失败")?;
if read == 0 {
break;
}
bytes += read as u64;
if bytes > 512 * 1024 * 1024 {
return Err("direct-execution-executor: 执行器读取超限".into());
}
digest.update(&buffer[..read]);
}
let after = file
.metadata()
.map_err(|_| "direct-execution-executor: 执行器身份不可用")?;
if bytes != before.len()
|| after.len() != before.len()
|| before.modified().ok() != after.modified().ok()
{
return Err("direct-execution-executor: 执行器在读取时改变".into());
}
Ok(format!("{:x}", digest.finalize()))
}
pub(super) struct ExecutionSession {
pub(super) root: PathBuf,
/// 本次确实新建执行账本;恢复和旧预算迁移均不构成新的用户受理。
pub(super) newly_accepted: bool,
state_path: PathBuf,
_owner: File,
data: Mutex<SessionData>,
analytics: Mutex<SessionAnalytics>,
changed: tokio::sync::watch::Sender<u64>,
cancellation: Arc<std::sync::atomic::AtomicBool>,
abort_requested: std::sync::atomic::AtomicBool,
}
pub(super) struct ExecutionSessionGuard {
pub(super) session: Arc<ExecutionSession>,
}
impl ExecutionSessionGuard {
pub(super) fn session(&self) -> Arc<ExecutionSession> {
Arc::clone(&self.session)
}
}
impl Drop for ExecutionSessionGuard {
fn drop(&mut self) {
if let Ok(mut sessions) = sessions().lock() {
if sessions
.get(&self.session.root)
.and_then(Weak::upgrade)
.is_some_and(|session| Arc::ptr_eq(&session, &self.session))
{
sessions.remove(&self.session.root);
}
}
if self.session.snapshot().is_ok_and(|state| {
!state.phase.is_terminal()
&& (state.requires_contract || state.contract.is_some() || !state.active.is_empty())
}) {
self.session
.cancellation
.store(true, std::sync::atomic::Ordering::Release);
self.session
.abort_requested
.store(true, std::sync::atomic::Ordering::Release);
let session = Arc::clone(&self.session);
let stop = move || {
let _ =
session.interrupt("回合已结束但宿主验收尚未完成,停止所有本轮执行。".into());
};
if let Ok(runtime) = tokio::runtime::Handle::try_current() {
runtime.spawn_blocking(stop);
} else {
stop();
}
}
}
}
pub(super) struct ExecutionLease {
session: Arc<ExecutionSession>,
id: String,
sequence: u32,
finished: bool,
}
#[derive(Clone)]
pub(crate) struct WritePermit {
session: Arc<ExecutionSession>,
id: String,
}
impl WritePermit {
pub(super) fn record_analytics_revision(
&self,
revision: u64,
change_kind: crate::analytics::contract::ChangeKind,
files_changed_count: u64,
) {
self.session
.record_analytics_revision(revision, change_kind, files_changed_count);
}
pub(crate) fn run<T>(&self, write: impl FnOnce() -> Result<T, String>) -> Result<T, String> {
let mut data = self.session.lock()?;
self.session.tick_locked(&mut data)?;
if data.ledger.phase.is_terminal() || data.ledger.phase == ExecutionPhase::Sealing {
return Err(closed_error(data.ledger.phase));
}
if self
.session
.cancellation
.load(std::sync::atomic::Ordering::Acquire)
|| self
.session
.abort_requested
.load(std::sync::atomic::Ordering::Acquire)
|| data
.ledger
.active
.get(&self.id)
.is_none_or(|record| record.kind != EffectKind::Write)
{
return Err("direct-execution-write-lease: 原写入许可已关闭,未提交项目修改".into());
}
// 只保护持有项目锁之后的同步本地提交,禁止把下载或其它网络等待放入此闭包。
let result = write();
drop(data);
result
}
}
impl ExecutionLease {
pub(super) fn write_permit(&self) -> Result<WritePermit, String> {
if self
.session
.lock()?
.ledger
.active
.get(&self.id)
.is_none_or(|record| record.kind != EffectKind::Write)
{
return Err("direct-execution-write-lease: 原写入许可已结束或类型不符".into());
}
Ok(WritePermit {
session: Arc::clone(&self.session),
id: self.id.clone(),
})
}
pub(super) fn paid_submission_scope(
&self,
) -> Result<super::direct_paid_submission::DirectPaidSubmissionScope, String> {
if self
.session
.lock()?
.ledger
.active
.get(&self.id)
.is_none_or(|record| record.kind != EffectKind::Paid)
{
return Err("direct-execution-paid-lease: 原付费许可已结束或类型不符".into());
}
let session = Arc::clone(&self.session);
let id = self.id.clone();
Ok(
super::direct_paid_submission::DirectPaidSubmissionScope::new(
Arc::new(move || session.begin_paid_dispatch(&id)),
self.session.cancel_flag(),
),
)
}
pub(super) fn sequence(&self) -> u32 {
self.sequence
}
pub(super) fn finish(
self,
passed: bool,
source_changed: bool,
evidence: Option<ExecutionEvidence>,
) -> Result<(), String> {
self.finish_with_dispatch_verdict(passed, source_changed, evidence, false)
}
pub(super) fn finish_paid_dispatch_denied(self, source_changed: bool) -> Result<(), String> {
self.finish_with_dispatch_verdict(false, source_changed, None, true)
}
fn finish_with_dispatch_verdict(
mut self,
passed: bool,
source_changed: bool,
evidence: Option<ExecutionEvidence>,
known_not_dispatched: bool,
) -> Result<(), String> {
let result = self.session.finish_lease(
&self.id,
passed,
source_changed,
evidence,
known_not_dispatched,
);
if result.is_err() {
self.session
.abort_requested
.store(true, std::sync::atomic::Ordering::Release);
let _ = self
.session
.interrupt("执行已返回,但宿主未能保存可信结算,已停止本轮。".into());
}
self.finished = true;
result
}
}
impl Drop for ExecutionLease {
fn drop(&mut self) {
if self.finished {
return;
}
self.session
.abort_requested
.store(true, std::sync::atomic::Ordering::Release);
self.session
.cancellation
.store(true, std::sync::atomic::Ordering::Release);
let session = Arc::clone(&self.session);
let id = self.id.clone();
let finish = move || {
let _ = session.finish_lease(&id, false, false, None, false);
let _ = session.interrupt("执行租约未正常结算,本轮已停止;请核对原执行结果。".into());
};
if let Ok(runtime) = tokio::runtime::Handle::try_current() {
runtime.spawn_blocking(finish);
} else {
finish();
}
}
}
fn sessions() -> &'static Mutex<BTreeMap<PathBuf, Weak<ExecutionSession>>> {
SESSIONS.get_or_init(|| Mutex::new(BTreeMap::new()))
}
fn now_ms() -> u64 {
SystemTime::now()
.duration_since(UNIX_EPOCH)
.unwrap_or_default()
.as_millis()
.min(u64::MAX as u128) as u64
}
fn hash(bytes: &[u8]) -> String {
format!("{:x}", Sha256::digest(bytes))
}
fn closed_error(phase: ExecutionPhase) -> String {
if phase == ExecutionPhase::Exhausted {
return "validation-budget-exhausted: 本轮预算已耗尽,禁止新的修改、执行和付费扩项".into();
}
format!("direct-execution-closed: 宿主执行状态为 {phase:?},请读取交付状态")
}
pub(super) fn current(root: &Path) -> Result<Arc<ExecutionSession>, String> {
let root = root
.canonicalize()
.map_err(|_| "direct-execution-project: 项目不可用")?;
sessions()
.lock()
.map_err(|_| "direct-execution-state: 会话状态不可用")?
.get(&root)
.and_then(Weak::upgrade)
.ok_or_else(|| "direct-execution-missing: 当前回合尚未建立宿主执行状态".into())
}
#[cfg(test)]
pub(super) fn register_for_test(
session: Arc<ExecutionSession>,
) -> Result<ExecutionSessionGuard, String> {
let mut registry = sessions().lock().map_err(|_| "执行测试注册表不可用")?;
if registry
.get(&session.root)
.and_then(Weak::upgrade)
.is_some()
{
return Err("测试项目已有执行会话".into());
}
registry.insert(session.root.clone(), Arc::downgrade(&session));
Ok(ExecutionSessionGuard { session })
}
pub(super) async fn begin(
root: &Path,
prompt: &str,
requires_contract: bool,
config: DirectValidationConfig,
analytics_run: Option<crate::analytics::run::Metadata>,
) -> Result<ExecutionSessionGuard, String> {
let root = root.to_path_buf();
let prompt_hash = hash(prompt.as_bytes());
let session = tokio::task::spawn_blocking(move || {
let turn = super::direct_taonier_active_invocation_id_at(&root)?;
let host = crate::game_creator_runtime_config_dir()
.ok_or("direct-execution-host: 需要客户端私有配置目录,CLI 请提供 --config-dir")?;
open_with_analytics_at(
&host.join("direct-executions"),
&root,
&turn,
&prompt_hash,
requires_contract,
&config,
analytics_run,
)
})
.await
.map_err(|_| "direct-execution-start: 宿主状态初始化退出")??;
let root = session.root.clone();
{
let mut registry = sessions()
.lock()
.map_err(|_| "direct-execution-state: 会话状态不可用")?;
if registry.get(&root).and_then(Weak::upgrade).is_some() {
return Err("direct-execution-active: 项目已有执行状态".into());
}
registry.insert(root, Arc::downgrade(&session));
}
let weak = Arc::downgrade(&session);
tokio::spawn(async move {
loop {
tokio::time::sleep(std::time::Duration::from_millis(250)).await;
let Some(session) = weak.upgrade() else {
break;
};
let done = tokio::task::spawn_blocking(move || {
if session.tick().is_err() {
return true;
}
session
.snapshot()
.map(|state| state.phase.is_terminal())
.unwrap_or(true)
})
.await
.unwrap_or(true);
if done {
break;
}
}
});
Ok(ExecutionSessionGuard { session })
}
#[cfg(test)]
pub(super) fn open_at(
host: &Path,
root: &Path,
turn: &str,
request_hash: &str,
requires_contract: bool,
config: &DirectValidationConfig,
) -> Result<Arc<ExecutionSession>, String> {
open_with_analytics_at(
host,
root,
turn,
request_hash,
requires_contract,
config,
None,
)
}
pub(super) fn open_with_analytics_at(
host: &Path,
root: &Path,
turn: &str,
request_hash: &str,
requires_contract: bool,
config: &DirectValidationConfig,
analytics_run: Option<crate::analytics::run::Metadata>,
) -> Result<Arc<ExecutionSession>, String> {
config.validate()?;
let root = root
.canonicalize()
.map_err(|_| "direct-execution-project: 项目不可用")?;
if turn.is_empty() || turn.len() > 512 || request_hash.len() != 64 {
return Err("direct-execution-identity: 回合身份无效".into());
}
if host.starts_with(&root) {
return Err("direct-execution-host: 权威状态不得写入模型项目目录".into());
}
crate::ensure_game_creator_private_directory_tree(host, "宿主执行状态目录")
.map_err(|_| "direct-execution-host: 无法准备私有执行状态目录")?;
let host = host
.canonicalize()
.map_err(|_| "direct-execution-host: 状态目录不可用")?;
if host.starts_with(&root) {
return Err("direct-execution-host: 权威状态不得写入模型项目目录".into());
}
let project_key = hash(root.to_string_lossy().as_bytes());
let directory = host.join(&project_key);
crate::ensure_game_creator_private_directory_tree(&directory, "宿主项目执行状态")
.map_err(|_| "direct-execution-host: 无法准备私有项目执行状态")?;
let key = hash(turn.as_bytes());
let lock_path = directory.join(format!("{key}.lock"));
if lock_path.exists() {
crate::prepare_game_creator_private_path_for_read(&lock_path, false, "执行归属锁")
.map_err(|_| "direct-execution-owner: 执行归属锁权限或类型无效")?;
}
let owner = OpenOptions::new()
.read(true)
.write(true)
.create(true)
.truncate(false)
.open(&lock_path)
.map_err(|_| "direct-execution-owner: 无法打开执行归属锁")?;
owner
.try_lock()
.map_err(|_| "direct-execution-active: 同一回合仍由其他执行器持有")?;
let state_path = directory.join(format!("{key}.json"));
let existing = if state_path.exists() {
Some(
serde_json::from_str::<ExecutionLedger>(
&crate::read_game_creator_private_file_to_string(
&state_path,
"宿主执行状态",
MAX_STATE_BYTES as u64,
)
.map_err(|_| {
"direct-execution-private-state: 无法安全读取宿主状态,禁止重置预算"
})?,
)
.map_err(|_| "direct-execution-corrupt: 宿主执行状态无效,禁止重置预算")?,
)
} else {
None
};
let project_id = super::read_existing_manifest_for_project(&root)?.project_id;
let is_new = existing.is_none();
let mut newly_accepted = is_new;
let mut ledger = existing.unwrap_or_else(|| ExecutionLedger {
schema_version: SCHEMA.into(),
client_turn_id: turn.into(),
project_id: project_id.clone(),
project_key: project_key.clone(),
request_hash: request_hash.into(),
requires_contract,
contract: None,
phase: ExecutionPhase::Working,
revision: 0,
used_passes: 0,
used_execution_ms: 0,
max_runs: config.max_runs,
max_execution_ms: config.max_execution_seconds.saturating_mul(1000),
max_turn_ms: config.max_turn_seconds.saturating_mul(1000),
created_at_ms: now_ms(),
next_sequence: 0,
validation_source: None,
active: BTreeMap::new(),
evidence: BTreeMap::new(),
executor_stopped: false,
terminal_report: None,
delivery_reviews: 0,
plan: None,
last_failed_write_revision: None,
analytics_run,
});
if is_new {
// 只继承旧项目账本的消费量,绝不把可编辑的旧成功回执提升为宿主证据。
let relative = format!(".agent/runtime/direct-validation/{key}.json");
let legacy: Option<Value> =
super::runtime_protocol::read_agent_runtime_json_sidecar_with_max_bytes(
&root,
&relative,
"旧验证预算",
512 * 1024,
)?;
if let Some(legacy) = legacy {
newly_accepted = false;
ledger.analytics_run = None;
let used = legacy["usedRuns"]
.as_u64()
.and_then(|n| u32::try_from(n).ok());
let maximum = legacy["maxRuns"]
.as_u64()
.and_then(|n| u32::try_from(n).ok())
.filter(|n| *n > 0);
if legacy["schemaVersion"] != "agc-direct-validation.v1"
|| legacy["clientTurnId"] != turn
|| used.is_none()
|| maximum.is_none()
{
return Err("direct-execution-legacy: 旧预算无效,不能重置消费量".into());
}
ledger.used_passes = used.unwrap();
ledger.max_runs = maximum.unwrap().min(config.max_runs);
ledger.next_sequence = ledger.used_passes;
if ledger.used_passes > 0 {
ledger.phase = ExecutionPhase::Draining;
}
if ledger.used_passes >= ledger.max_runs {
ledger.phase = ExecutionPhase::Exhausted;
ledger.terminal_report =
Some("旧回合的验证预算已耗尽,不能因启用宿主控制而重新开始执行。".into());
}
}
}
if ledger.schema_version != SCHEMA
|| ledger.client_turn_id != turn
|| ledger.project_key != project_key
|| ledger.project_id != project_id
|| ledger.request_hash != request_hash
|| ledger.max_runs == 0
|| ledger.max_execution_ms == 0
|| ledger.max_turn_ms == 0
{
return Err("direct-execution-identity: 持久状态与当前回合不一致,禁止重置预算".into());
}
if (!ledger.active.is_empty() || ledger.phase == ExecutionPhase::Sealing)
&& !ledger.phase.is_terminal()
{
ledger.phase = ExecutionPhase::Interrupted;
ledger.terminal_report = Some(
"上一执行器退出时仍有操作未结算。本轮保持未完成,需核对原操作结果,不能自动重放。"
.into(),
);
}
if now_ms().saturating_add(1000) < ledger.created_at_ms {
ledger.phase = ExecutionPhase::Interrupted;
ledger.terminal_report = Some("宿主时钟发生回退,无法证明原执行期限,已停止本轮。".into());
}
ledger.requires_contract |= requires_contract;
let initial_elapsed_ms = now_ms().saturating_sub(ledger.created_at_ms);
let (changed, _) = tokio::sync::watch::channel(ledger.revision);
let session = Arc::new(ExecutionSession {
root,
newly_accepted,
state_path,
_owner: owner,
analytics: Mutex::new(SessionAnalytics {
project_id: ledger.project_id.clone(),
route: ledger
.analytics_run
.as_ref()
.map(|run| run.context.route.clone()),
capture: None,
output_revision: None,
}),
data: Mutex::new(SessionData {
ledger,
started: Instant::now(),
initial_elapsed_ms,
lease_started: BTreeMap::new(),
poisoned: false,
codex_executor: None,
#[cfg(test)]
elapsed_offset_ms: 0,
}),
changed,
cancellation: Arc::new(std::sync::atomic::AtomicBool::new(false)),
abort_requested: std::sync::atomic::AtomicBool::new(false),
});
{
let mut data = session.lock()?;
let next = data.ledger.clone();
session.commit(&mut data, next)?;
}
session.tick()?;
Ok(session)
}
impl ExecutionSession {
pub(super) fn bind_codex_executor(&self, path: &Path, version: &str) -> Result<(), String> {
// 发行构建只接受捆绑侧车固定版本;开发构建用宿主自带的 Codex,按 profile 跳过该门禁。
#[cfg(not(debug_assertions))]
{
if version.trim() != super::codex_cli::codex_bundle::CLI_VERSION {
return Err(format!(
"direct-execution-executor: 尚未验证该执行器的补丁协议(期望 {},实际 {})",
super::codex_cli::codex_bundle::CLI_VERSION,
version.trim()
));
}
}
#[cfg(debug_assertions)]
let _ = version;
let path = path
.canonicalize()
.map_err(|_| "direct-execution-executor: 无法锚定执行器")?;
let identity = CodexExecutorIdentity {
digest: executor_digest(&path)?,
path,
version: version.trim().into(),
};
let mut data = self.lock()?;
if data
.codex_executor
.as_ref()
.is_some_and(|existing| existing != &identity)
{
return Err("direct-execution-executor: 同一回合不能切换执行器".into());
}
data.codex_executor = Some(identity);
Ok(())
}
pub(super) fn codex_executor(&self) -> Result<PathBuf, String> {
let identity = self
.lock()?
.codex_executor
.clone()
.ok_or("direct-execution-executor: 当前回合未绑定补丁执行器")?;
if executor_digest(&identity.path)? != identity.digest {
return Err("direct-execution-executor: 已绑定执行器发生变化,不能执行补丁".into());
}
Ok(identity.path)
}
fn lock(&self) -> Result<std::sync::MutexGuard<'_, SessionData>, String> {
self.data
.lock()
.map_err(|_| "direct-execution-state: 执行状态不可用".into())
}
fn commit(&self, data: &mut SessionData, mut next: ExecutionLedger) -> Result<(), String> {
if data.poisoned {
return Err("direct-execution-persistence: 状态落盘失败,已关闭执行".into());
}
next.revision = data.ledger.revision.saturating_add(1);
let bytes = serde_json::to_vec(&next)
.map_err(|_| "direct-execution-persistence: 状态序列化失败")?;
if bytes.len() > MAX_STATE_BYTES {
return Err("direct-execution-capacity: 宿主状态超过保留上限".into());
}
if crate::write_game_creator_private_file(&self.state_path, &bytes, "宿主执行状态").is_err()
{
data.poisoned = true;
data.ledger.phase = ExecutionPhase::Interrupted;
data.ledger.terminal_report = Some("执行状态未能持久化,已停止接受新操作。".into());
self.cancellation
.store(true, std::sync::atomic::Ordering::Release);
self.changed.send_replace(next.revision);
return Err("direct-execution-persistence: 状态落盘失败,已关闭执行".into());
}
data.ledger = next;
if data.ledger.phase.is_terminal() || data.ledger.phase == ExecutionPhase::Sealing {
self.cancellation
.store(true, std::sync::atomic::Ordering::Release);
}
self.changed.send_replace(data.ledger.revision);
Ok(())
}
fn elapsed(data: &SessionData) -> u64 {
let elapsed = data
.initial_elapsed_ms
.saturating_add(data.started.elapsed().as_millis().min(u64::MAX as u128) as u64);
#[cfg(test)]
let elapsed = elapsed.saturating_add(data.elapsed_offset_ms);
elapsed
}
fn running_ms(data: &SessionData) -> u64 {
data.ledger
.active
.iter()
.filter(|(_, entry)| entry.kind != EffectKind::Write)
.map(|(id, _)| {
let duration = data
.lease_started
.get(id)
.map(|at| at.elapsed().as_millis().min(u64::MAX as u128) as u64)
.unwrap_or(0);
#[cfg(test)]
let duration = duration.saturating_add(data.elapsed_offset_ms);
duration
})
.fold(0u64, u64::saturating_add)
}
pub(super) fn subscribe(&self) -> tokio::sync::watch::Receiver<u64> {
self.changed.subscribe()
}
pub(super) fn cancel_flag(&self) -> Arc<std::sync::atomic::AtomicBool> {
Arc::clone(&self.cancellation)
}
pub(super) fn was_aborted(&self) -> bool {
self.abort_requested
.load(std::sync::atomic::Ordering::Acquire)
}
pub(super) fn record_delivery_review(&self) -> Result<u32, String> {
let mut data = self.lock()?;
if data.ledger.phase.is_terminal() {
return Ok(data.ledger.delivery_reviews);
}
let mut next = data.ledger.clone();
next.delivery_reviews = next.delivery_reviews.saturating_add(1);
let count = next.delivery_reviews;
self.commit(&mut data, next)?;
Ok(count)
}
pub(super) fn update_plan(&self, plan: Value) -> Result<Value, String> {
self.tick()?;
if !plan.is_object()
|| serde_json::to_vec(&plan).map_err(|_| "计划格式无效")?.len() > 32 * 1024
{
return Err("direct-plan-invalid: 计划必须为有界对象".into());
}
let mut data = self.lock()?;
if data.ledger.phase.is_terminal() || data.ledger.phase == ExecutionPhase::Sealing {
return Err(closed_error(data.ledger.phase));
}
let mut next = data.ledger.clone();
next.plan = Some(plan.clone());
self.commit(&mut data, next)?;
Ok(json!({"plan":plan,"revision":data.ledger.revision,"acceptancePassed":false}))
}
pub(super) fn set_analytics_capture(
&self,
capture: Option<(
crate::analytics::contract::Context,
crate::analytics::store::AnalyticsWriter,
)>,
) {
let Ok(mut analytics) = self.analytics.lock() else {
return;
};
analytics.capture = capture.and_then(|(mut context, writer)| {
// 恢复或账号切换后仍归属于真实受理的原 run。
context.route = analytics.route.clone()?;
Some((context, writer))
});
}
pub(super) fn analytics_capture(
&self,
) -> Option<(
crate::analytics::contract::Context,
crate::analytics::store::AnalyticsWriter,
)> {
self.analytics.lock().ok()?.capture.clone()
}
pub(super) fn record_analytics_revision(
&self,
revision: u64,
change_kind: crate::analytics::contract::ChangeKind,
files_changed_count: u64,
) {
use crate::analytics::contract::{RevisionCreated, RevisionSource, Source};
if files_changed_count == 0 {
return;
}
let Ok(mut analytics) = self.analytics.lock() else {
return;
};
if analytics.route.is_none() {
return;
}
analytics.output_revision = Some(analytics.output_revision.unwrap_or(0).max(revision));
let capture = analytics.capture.clone();
let project_id = analytics.project_id.clone();
drop(analytics);
crate::analytics::project::revision(
capture,
&project_id,
Source::Direct,
RevisionCreated {
revision_id: revision.to_string(),
revision_source: RevisionSource::Agent,
change_kind,
files_changed_count: Some(files_changed_count),
},
);
}
pub(super) fn analytics_output_revision(&self) -> Option<String> {
self.analytics
.lock()
.ok()?
.output_revision
.map(|revision| revision.to_string())
}
pub(super) fn snapshot(&self) -> Result<ExecutionLedger, String> {
let data = self.lock()?;
let mut state = data.ledger.clone();
state.used_execution_ms = state
.used_execution_ms
.saturating_add(Self::running_ms(&data));
Ok(state)
}
pub(super) fn tick(&self) -> Result<(), String> {
let mut data = self.lock()?;
self.tick_locked(&mut data)
}
fn tick_locked(&self, data: &mut SessionData) -> Result<(), String> {
if data.ledger.phase.is_terminal() {
return Ok(());
}
let execution_ms = data
.ledger
.used_execution_ms
.saturating_add(Self::running_ms(&data));
let elapsed = Self::elapsed(&data);
if execution_ms >= data.ledger.max_execution_ms || elapsed >= data.ledger.max_turn_ms {
let mut next = data.ledger.clone();
next.phase = ExecutionPhase::Exhausted;
next.terminal_report = Some(format!("本轮预算已耗尽,交付尚未完成。已用执行批次 {}/{};累计工具执行约 {} 秒,整轮耗时约 {} 秒。保留已有证据与未完成项,停止新的修改、执行和付费扩项。", next.used_passes, next.max_runs, execution_ms / 1000, elapsed / 1000));
self.commit(data, next)?;
}
Ok(())
}
fn begin_paid_dispatch(&self, id: &str) -> Result<(), String> {
let mut data = self.lock()?;
self.tick_locked(&mut data)?;
if data.ledger.phase.is_terminal() || data.ledger.phase == ExecutionPhase::Sealing {
return Err(closed_error(data.ledger.phase));
}
if self.cancellation.load(std::sync::atomic::Ordering::Acquire)
|| self
.abort_requested
.load(std::sync::atomic::Ordering::Acquire)
{
return Err("direct-execution-interrupted: 付费提交前已取消".into());
}
let mut next = data.ledger.clone();
let record = next
.active
.get_mut(id)
.filter(|record| record.kind == EffectKind::Paid)
.ok_or("direct-execution-paid-lease: 原付费许可已结束或类型不符")?;
record.paid_dispatch_count = record
.paid_dispatch_count
.checked_add(1)
.ok_or("direct-execution-paid-lease: 付费提交序号耗尽")?;
// 与 begin_sealing 使用同一状态锁。通过后只允许此请求完成/对账,下一次 POST 仍须复核。
self.commit(&mut data, next)
}
pub(super) fn freeze_contract(&self, contract: Value) -> Result<Value, String> {
if !contract.is_object()
|| serde_json::to_vec(&contract)
.map_err(|_| "交付合同无效")?
.len()
> 64 * 1024
{
return Err("direct-execution-contract: 合同必须为有界对象".into());
}
self.tick()?;
let mut data = self.lock()?;
if let Some(existing) = &data.ledger.contract {
return if *existing == contract {
Ok(existing.clone())
} else {
Err(
"direct-execution-contract-frozen: 本轮合同已冻结,新增范围需要新的用户回合"
.into(),
)
};
}
if !matches!(
data.ledger.phase,
ExecutionPhase::Working | ExecutionPhase::Draining
) {
return Err(closed_error(data.ledger.phase));
}
let mut next = data.ledger.clone();
next.contract = Some(contract.clone());
next.requires_contract = true;
self.commit(&mut data, next)?;
Ok(contract)
}
pub(super) fn admit(
self: &Arc<Self>,
kind: EffectKind,
validation_source: Option<String>,
) -> Result<ExecutionLease, String> {
self.admit_with_key(kind, validation_source, None)
}
pub(super) fn admit_validation(
self: &Arc<Self>,
key: &str,
source: &str,
) -> Result<ExecutionLease, String> {
self.admit_with_key(EffectKind::Execute, Some(source.into()), Some(key.into()))
}
fn admit_with_key(
self: &Arc<Self>,
kind: EffectKind,
validation_source: Option<String>,
validation_key: Option<String>,
) -> Result<ExecutionLease, String> {
if self
.abort_requested
.load(std::sync::atomic::Ordering::Acquire)
{
return Err("direct-execution-interrupted: 原执行未正常结算,不能接受新操作".into());
}
self.tick()?;
let mut data = self.lock()?;
let state = &data.ledger;
if state.phase.is_terminal() || state.phase == ExecutionPhase::Sealing {
return Err(closed_error(state.phase));
}
if state.contract.is_none() {
return Err("direct-execution-contract-required: 先登记本轮必需范围与验收合同,再执行修改或命令".into());
}
if state.phase == ExecutionPhase::Draining && state.used_passes >= state.max_runs {
let mut next = state.clone();
next.phase = ExecutionPhase::Exhausted;
next.terminal_report =
Some("本轮返修预算已耗尽,停止新的修改、执行和付费扩项。".into());
self.commit(&mut data, next)?;
return Err(closed_error(ExecutionPhase::Exhausted));
}
if state.active.len() >= MAX_ACTIVE {
return Err("direct-execution-capacity: 请等待已受理操作结束".into());
}
if validation_key.is_some()
&& state.active.values().any(|entry| {
entry.validation_key == validation_key
&& entry.validation_source == validation_source
})
{
return Err("validation-already-running: 同一输入的验证已受理,请等待原执行".into());
}
let mut next = state.clone();
if validation_source
.as_ref()
.zip(next.validation_source.as_ref())
.is_some_and(|(a, b)| a != b)
{
next.phase = ExecutionPhase::Draining;
}
if next.phase == ExecutionPhase::Draining && !next.active.is_empty() {
if data.ledger.phase != next.phase {
self.commit(&mut data, next)?;
}
return Err("direct-execution-draining: 当前批次正在收束,请等待在途操作结束".into());
}
if kind != EffectKind::Write
&& (next.used_passes == 0 || next.phase == ExecutionPhase::Draining)
{
if next.used_passes >= next.max_runs {
next.phase = ExecutionPhase::Exhausted;
next.terminal_report = Some(format!(
"本轮执行/返修批次已耗尽({}/{}),交付尚未完成。保留已有证据并停止扩项。",
next.used_passes, next.max_runs
));
self.commit(&mut data, next)?;
return Err("validation-budget-exhausted: 执行/返修批次已耗尽".into());
}
next.used_passes += 1;
next.phase = ExecutionPhase::Working;
next.validation_source = None;
}
if let Some(source) = &validation_source {
next.validation_source = Some(source.clone());
}
next.next_sequence = next
.next_sequence
.checked_add(1)
.ok_or("direct-execution-sequence: 执行序号耗尽")?;
let sequence = next.next_sequence;
let id = uuid::Uuid::new_v4().to_string();
next.executor_stopped = false;
next.active.insert(
id.clone(),
LeaseRecord {
kind,
sequence,
pass: next.used_passes,
started_at_ms: now_ms(),
validation_source,
validation_key,
admitted_revision: data.ledger.revision.saturating_add(1),
paid_dispatch_count: 0,
},
);
self.commit(&mut data, next)?;
data.lease_started.insert(id.clone(), Instant::now());
self.cancellation
.store(false, std::sync::atomic::Ordering::Release);
Ok(ExecutionLease {
session: Arc::clone(self),
id,
sequence,
finished: false,
})
}
fn finish_lease(
&self,
id: &str,
passed: bool,
source_changed: bool,
evidence: Option<ExecutionEvidence>,
known_not_dispatched: bool,
) -> Result<(), String> {
let mut data = self.lock()?;
let mut next = data.ledger.clone();
let Some(record) = next.active.remove(id) else {
return Err("direct-execution-receipt: 租约不存在或已经完成".into());
};
if record.kind == EffectKind::Write {
if !passed {
next.last_failed_write_revision = Some(data.ledger.revision.saturating_add(1));
} else if next
.last_failed_write_revision
.is_some_and(|failed| record.admitted_revision > failed)
{
next.last_failed_write_revision = None;
}
}
if record.kind != EffectKind::Write {
next.used_execution_ms = next.used_execution_ms.saturating_add(
data.lease_started
.get(id)
.map(|at| at.elapsed().as_millis().min(u64::MAX as u128) as u64)
.unwrap_or(0),
);
}
// 封口主动回收辅助 native 会话只结算租约,不把取消伪造为验证成功;
// 可信验证失败仍写入 evidence,最终复核必须重新读取该证据与实际文件。
if !next.phase.is_terminal()
&& next.phase != ExecutionPhase::Sealing
&& ((!passed && record.kind != EffectKind::Write)
|| (source_changed && next.validation_source.is_some()))
{
next.phase = ExecutionPhase::Draining;
}
if next.phase == ExecutionPhase::Draining && next.used_passes >= next.max_runs {
next.phase = ExecutionPhase::Exhausted;
next.terminal_report=Some(format!("本轮执行/返修批次已耗尽({}/{}),最近一次验证仍未通过,停止新的修改、执行和付费扩项。",next.used_passes,next.max_runs));
}
if !passed
&& record.kind == EffectKind::Paid
&& !(known_not_dispatched && record.paid_dispatch_count == 0)
&& next.phase == ExecutionPhase::Sealing
{
next.phase = ExecutionPhase::Interrupted;
next.terminal_report = Some(
"收尾时仍有未成功结算的付费操作,保留原操作记录并停止本轮,不能自动重放。".into(),
);
}
if let Some(mut evidence) = evidence {
if !evidence.result.is_object() {
return Err("direct-execution-evidence: 可信验证结果必须为对象".into());
}
if evidence.key.len() > 256
|| (!next.evidence.contains_key(&evidence.key)
&& next.evidence.len() >= MAX_EVIDENCE)
{
return Err("direct-execution-evidence-capacity: 验证证据数量超限".into());
}
let unresolved = evidence.result["needsReconciliation"] == true
|| evidence.result["timedOut"] == true;
evidence.result["passed"] = json!(passed && !source_changed && !unresolved);
evidence.result["sourceChanged"] = json!(source_changed);
if passed
&& !source_changed
&& !unresolved
&& next
.last_failed_write_revision
.is_some_and(|failed| record.admitted_revision > failed)
{
next.last_failed_write_revision = None;
}
if unresolved {
next.phase = ExecutionPhase::Interrupted;
let prior = next.terminal_report.take().unwrap_or_default();
next.terminal_report = Some(format!(
"{prior}\n执行超时或执行结果未确认,本轮已停止;请核对原操作,不能自动重试。"
));
}
next.evidence.insert(evidence.key.clone(), evidence);
}
self.commit(&mut data, next)?;
data.lease_started.remove(id);
drop(data);
self.tick()
}
pub(super) fn begin_sealing(&self, expected_revision: u64) -> Result<bool, String> {
self.tick()?;
let mut data = self.lock()?;
if data.ledger.revision != expected_revision
|| data.ledger.phase != ExecutionPhase::Working
|| data.ledger.contract.is_none()
|| data.ledger.last_failed_write_revision.is_some()
|| data
.ledger
.active
.values()
.any(|entry| entry.kind == EffectKind::Write)
{
return Ok(false);
}
let mut next = data.ledger.clone();
next.phase = ExecutionPhase::Sealing;
self.commit(&mut data, next)?;
Ok(true)
}
pub(super) fn record_process_exit_proof(&self, proven: bool) -> Result<(), String> {
let mut data = self.lock()?;
if data.ledger.phase == ExecutionPhase::Completed && !proven {
return Err("direct-execution-proof: 不能覆盖已完成回合的退出证明".into());
}
let mut next = data.ledger.clone();
next.executor_stopped = proven;
if !proven && next.phase != ExecutionPhase::Completed {
next.phase = ExecutionPhase::Interrupted;
let prior = next.terminal_report.take().unwrap_or_default();
next.terminal_report = Some(format!(
"{prior}\n执行器缺少完整子进程退出证明。本轮仍未验收,需核对后台操作结果。"
));
}
self.commit(&mut data, next)
}
pub(super) fn complete(&self, report: String) -> Result<(), String> {
let mut data = self.lock()?;
if data.ledger.phase != ExecutionPhase::Sealing
|| !data.ledger.active.is_empty()
|| !data.ledger.executor_stopped
{
return Err("direct-execution-completion: 缺少封口、排空或执行器退出证明".into());
}
let mut next = data.ledger.clone();
next.phase = ExecutionPhase::Completed;
next.terminal_report = Some(report);
self.commit(&mut data, next)
}
pub(super) fn reopen_for_repair(&self) -> Result<(), String> {
let mut data = self.lock()?;
if data.ledger.phase != ExecutionPhase::Sealing {
return Err(closed_error(data.ledger.phase));
}
let mut next = data.ledger.clone();
next.phase = if next.used_passes >= next.max_runs {
ExecutionPhase::Exhausted
} else {
ExecutionPhase::Draining
};
if next.phase == ExecutionPhase::Exhausted {
next.terminal_report =
Some("最终复核未通过且本轮返修预算已耗尽,保留现有修改与证据,停止扩项。".into());
}
next.executor_stopped = false;
self.commit(&mut data, next)
}
pub(super) fn interrupt(&self, reason: String) -> Result<(), String> {
let mut data = self.lock()?;
if matches!(
data.ledger.phase,
ExecutionPhase::Completed | ExecutionPhase::Interrupted
) {
return Ok(());
}
let mut next = data.ledger.clone();
next.phase = ExecutionPhase::Interrupted;
let prior = next.terminal_report.take().unwrap_or_default();
next.terminal_report = Some(
format!("{prior}\n{reason}")
.trim()
.chars()
.take(16_000)
.collect(),
);
self.commit(&mut data, next)
}
}
#[cfg(test)]
mod tests;