485ed50b26
Project CI / AI game creator shell Rust smoke (push) Has been cancelled
Project CI / AI game creator shell Rust crates (push) Has been cancelled
Project CI / Backend tests (push) Has been cancelled
Project CI / Native shell tests (push) Has been cancelled
Project CI / Frontend tests (push) Has been cancelled
Project CI / Repository checks (push) Has been cancelled
Project CI / AI game creator shell web tests (push) Has been cancelled
Project CI / AI game creator shell Rust shard 2/4 (push) Has been cancelled
Project CI / AI game creator shell Rust shard 1/4 (push) Has been cancelled
Project CI / AI game creator shell Rust shard 3/4 (push) Has been cancelled
Project CI / AI game creator shell Rust shard 4/4 (push) Has been cancelled
## 背景 净月潭案例(2026-09-19 21:53 → 09-20 01:26,3 小时 33 分)的问题不是模型慢,而是环境未就绪、验收靠模型自述、预算只约束单个工具入口、工具调用被 SDK 全局串行闸门卡住。本 PR 落地确认后的七项交付效率合同,并把捆绑 Codex 升到当前 npm latest。 ## 结果 - 新建 Web 游戏在正式生成前由宿主自动预检:捆绑 Node/npm、真实 Vite 构建、受限浏览器桌面/移动截图;失败不启动生成或付费素材。 - 首个副作用前冻结交付合同,宿主保存权威证据与预算;证据齐全后先封口、排空并取得执行器完整退出证明,再产出交付报告,模型回复不再当作验收。 - 原生 shell、内置浏览器与托管命令、第三方 MCP 共用同一执行许可和累计执行时间;`validation.maxRuns` 改为执行/返修批次语义,另设 `maxTurnSeconds` 墙钟上限。 - 所有工具可并发:移除 SDK 全局串行的 `apply_patch`/`update_plan` 注册,改由宿主 `agc_apply_patch`(官方 parser + 当前回合写许可 + 受控进程树 + 短项目事务)和 `agc_update_plan` 提供等价能力;真实请求目录里已无串行注册,长 MCP 与补丁/计划实测在同一响应内重叠。 - 付费提交与本地写入绑定原回合原租约:容量与同动作锁等待可取消,封口、取消或预算耗尽后零新增提交;已越过提交边界的请求保留 operation ID 走 GET 对账,不自动重放。 - 新增请求与工具分段计时、有界并行批读和首轮上下文预取,未知耗时不补零。 - 捆绑 Codex 0.147.0 → 0.155.1:固定版本收敛到 `build_support/codex_bundle.rs` 单一声明,vendor 解析源码按 `rust-v0.155.1` 逐字节重取并更新 UPSTREAM 证据,适配 0.155 统一 exec(`exec_command`/`write_stdin`);macOS 侧车最小系统版本仍为 15.0。 ## 验证 - 生产 CLI → 捆绑 Codex 0.155.1 → 本地 Responses/MCP 夹具 9/9:补丁、完成收尾、批次、只读/可写 MCP 并发、原生命令并发、原生资源并发、统一 exec 会话、期限终止(真实重叠 775 / 764 / 999 ms)。 - 真实模型目录 2/2;vendor 上游库 111 项;宿主定向与回归 106 项;真实 Windows 进程与取消 21 项;Node 侧门禁 52 项。 - 发行载荷:artifact-only 重建 NSIS,解包验证侧车清单 `codex-cli 0.155.1`、6 个组件哈希、打包后 `bin/codex.exe --version`、Node 运行时 2130 文件与双端 PNG、`--environment-check` ready。 - `check-config`、TypeScript、编码(4991 文件)、文档索引、`git diff --check`、`cargo fmt --check` 全部通过。 ## 未覆盖 - 真实陶泥儿登录态 Provider 生成尚未执行(本机 `--llm-status` 返回 `authentication-required`)。 - macOS 侧车只在本机做静态 Mach-O 与清单解析,未在 macOS 上跑 `check-macos-bundle.mjs`。 - 全仓库聚合套件在本机仍因负载敏感的后台 mock-provider 用例而红,与本次改动无关:改动前的旧二进制同样失败,换单线程后相关用例 3/3 通过。 --------- Co-authored-by: kdletters <61648117+kdletters@users.noreply.github.com> Reviewed-on: http://192.168.35.82/git/GenarrativeAI/Genarrative/pulls/439
631 lines
23 KiB
Rust
631 lines
23 KiB
Rust
use super::*;
|
|
|
|
pub(super) const PROCESS_SESSION_SCHEMA_VERSION: &str = "3";
|
|
pub(super) const PROCESS_SESSION_TRANSCRIPT_SCHEMA_VERSION: &str = "2";
|
|
pub(super) const PROCESS_SESSION_CURSOR_VERSION: &str = "v1";
|
|
pub(super) const PROCESS_SESSION_MAX_PER_PROJECT: usize = 4;
|
|
pub(super) const PROCESS_SESSION_MAX_PER_AGENT: usize = 2;
|
|
pub(super) const PROCESS_SESSION_MAX_OUTPUT_BYTES: usize = 256 * 1024;
|
|
pub(super) const PROCESS_SESSION_MAX_PENDING_LINE_BYTES: usize = 16 * 1024;
|
|
pub(super) const PROCESS_SESSION_MAX_STDIN_BYTES: usize = 8 * 1024;
|
|
pub(super) const PROCESS_SESSION_DEFAULT_POLL_CHARS: usize = 8_000;
|
|
pub(super) const PROCESS_SESSION_MAX_POLL_CHARS: usize = 16_000;
|
|
pub(super) const PROCESS_SESSION_MAX_POLL_WAIT_MS: u64 = 30_000;
|
|
pub(super) const PROCESS_SESSION_RECORD_MAX_BYTES: usize = 32 * 1024;
|
|
pub(super) const PROCESS_SESSION_TRANSCRIPT_MAX_BYTES: usize = 320 * 1024;
|
|
pub(super) const PROCESS_SESSION_TERMINATE_GRACE_MS: u64 = 800;
|
|
#[cfg(target_os = "linux")]
|
|
pub(super) const PROCESS_SESSION_OWNER_PID_ENV: &str = "GENARRATIVE_PROCESS_SESSION_OWNER_PID";
|
|
#[cfg(target_os = "linux")]
|
|
pub(super) const PROCESS_SESSION_CHILD_MODE: &str = "--process-session-child";
|
|
|
|
#[derive(Clone, Debug, Eq, PartialEq)]
|
|
pub(crate) struct ProcessSessionIdentity {
|
|
pub(crate) project_id: String,
|
|
pub(crate) agent_id: String,
|
|
pub(crate) task_id: String,
|
|
pub(crate) conversation_session_id: String,
|
|
pub(crate) run_id: String,
|
|
pub(crate) start_action_id: String,
|
|
pub(crate) start_action_fingerprint: String,
|
|
}
|
|
|
|
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
|
|
#[serde(deny_unknown_fields, rename_all = "camelCase")]
|
|
pub(crate) struct ProcessSessionRecord {
|
|
pub(crate) schema_version: String,
|
|
pub(crate) project_id: String,
|
|
pub(crate) agent_id: String,
|
|
pub(crate) task_id: String,
|
|
pub(crate) conversation_session_id: String,
|
|
pub(crate) run_id: String,
|
|
pub(crate) start_action_id: String,
|
|
pub(crate) start_action_fingerprint: String,
|
|
pub(crate) process_id: String,
|
|
pub(crate) owner_boot_id: String,
|
|
pub(crate) command_id: String,
|
|
pub(crate) program: String,
|
|
pub(crate) cwd: String,
|
|
#[serde(default)]
|
|
pub(crate) sandbox_backend: String,
|
|
#[serde(default)]
|
|
pub(crate) sandbox_mode: String,
|
|
#[serde(default)]
|
|
pub(crate) network_access: String,
|
|
#[serde(default)]
|
|
pub(crate) sandbox_profile_version: String,
|
|
#[serde(default)]
|
|
pub(crate) sandbox_establishment: String,
|
|
#[serde(default)]
|
|
pub(crate) target_exec: String,
|
|
#[serde(default)]
|
|
pub(crate) launch_failure_kind: Option<String>,
|
|
#[serde(default)]
|
|
pub(crate) sandbox_ready_at: Option<u64>,
|
|
#[serde(default)]
|
|
pub(crate) exec_established_at: Option<u64>,
|
|
pub(crate) status: String,
|
|
pub(crate) exit_code: Option<i32>,
|
|
pub(crate) signal: Option<String>,
|
|
pub(crate) stdin_open: bool,
|
|
pub(crate) output_bytes: usize,
|
|
pub(crate) output_sha256: String,
|
|
pub(crate) output_ref: Option<String>,
|
|
pub(crate) source_fingerprint_before: String,
|
|
pub(crate) source_fingerprint_after: Option<String>,
|
|
pub(crate) source_changed: Option<bool>,
|
|
pub(crate) needs_reconciliation: bool,
|
|
pub(crate) started_at: u64,
|
|
pub(crate) terminal_at: Option<u64>,
|
|
pub(crate) updated_at: u64,
|
|
}
|
|
|
|
#[derive(Clone, Debug, Deserialize, Eq, PartialEq, Serialize)]
|
|
#[serde(deny_unknown_fields, rename_all = "camelCase")]
|
|
pub(super) struct ProcessSessionTranscript {
|
|
pub(super) schema_version: String,
|
|
pub(super) project_id: String,
|
|
pub(super) agent_id: String,
|
|
pub(super) task_id: String,
|
|
pub(super) conversation_session_id: String,
|
|
pub(super) run_id: String,
|
|
pub(super) start_action_id: String,
|
|
pub(super) start_action_fingerprint: String,
|
|
pub(super) process_id: String,
|
|
pub(super) output: String,
|
|
pub(super) output_sha256: String,
|
|
pub(super) output_bytes: usize,
|
|
pub(super) updated_at: u64,
|
|
}
|
|
|
|
#[derive(Clone, Debug, Eq, PartialEq, Serialize)]
|
|
#[serde(rename_all = "camelCase")]
|
|
pub(crate) struct ProcessSessionPollResult {
|
|
pub(crate) process_id: String,
|
|
pub(crate) status: String,
|
|
pub(crate) output: String,
|
|
pub(crate) cursor: String,
|
|
pub(crate) next_cursor: String,
|
|
pub(crate) has_more: bool,
|
|
pub(crate) stdin_open: bool,
|
|
pub(crate) exit_code: Option<i32>,
|
|
pub(crate) signal: Option<String>,
|
|
pub(crate) output_bytes: usize,
|
|
pub(crate) output_sha256: String,
|
|
pub(crate) source_changed: Option<bool>,
|
|
pub(crate) needs_reconciliation: bool,
|
|
pub(crate) sandbox_backend: String,
|
|
pub(crate) sandbox_mode: String,
|
|
pub(crate) network_access: String,
|
|
pub(crate) sandbox_profile_version: String,
|
|
pub(crate) sandbox_establishment: String,
|
|
pub(crate) target_exec: String,
|
|
pub(crate) launch_failure_kind: Option<String>,
|
|
}
|
|
|
|
#[derive(Clone, Debug, Eq, PartialEq, Serialize)]
|
|
#[serde(rename_all = "camelCase")]
|
|
pub(crate) struct ProcessSessionStdinResult {
|
|
pub(crate) process_id: String,
|
|
pub(crate) bytes_written: usize,
|
|
pub(crate) content_sha256: String,
|
|
pub(crate) stdin_open: bool,
|
|
pub(crate) eof: bool,
|
|
pub(crate) sandbox_backend: String,
|
|
pub(crate) sandbox_mode: String,
|
|
pub(crate) network_access: String,
|
|
pub(crate) sandbox_profile_version: String,
|
|
}
|
|
|
|
#[derive(Debug)]
|
|
pub(super) struct ProcessOutputState {
|
|
pub(super) text: String,
|
|
pub(super) status: String,
|
|
pub(super) exit_code: Option<i32>,
|
|
pub(super) signal: Option<String>,
|
|
pub(super) stdin_open: bool,
|
|
pub(super) reader_finished: bool,
|
|
pub(super) output_limit_exceeded: bool,
|
|
pub(super) source_fingerprint_after: Option<String>,
|
|
pub(super) source_changed: Option<bool>,
|
|
pub(super) needs_reconciliation: bool,
|
|
pub(super) launch_failure_kind: Option<String>,
|
|
}
|
|
|
|
impl ProcessOutputState {
|
|
pub(super) fn running() -> Self {
|
|
Self {
|
|
text: String::new(),
|
|
status: "running".to_string(),
|
|
exit_code: None,
|
|
signal: None,
|
|
stdin_open: true,
|
|
reader_finished: false,
|
|
output_limit_exceeded: false,
|
|
source_fingerprint_after: None,
|
|
source_changed: None,
|
|
needs_reconciliation: false,
|
|
launch_failure_kind: None,
|
|
}
|
|
}
|
|
}
|
|
|
|
#[derive(Debug)]
|
|
pub(super) enum ProcessControl {
|
|
Terminate,
|
|
OutputLimit,
|
|
Shutdown,
|
|
}
|
|
|
|
pub(super) struct LiveProcessSession {
|
|
pub(super) root: PathBuf,
|
|
pub(super) identity: ProcessSessionIdentity,
|
|
pub(super) process_id: String,
|
|
pub(super) command_id: String,
|
|
pub(super) program: String,
|
|
pub(super) cwd: String,
|
|
pub(super) sandbox_backend: String,
|
|
pub(super) sandbox_mode: String,
|
|
pub(super) network_access: String,
|
|
pub(super) sandbox_profile_version: String,
|
|
pub(super) sandbox_establishment: String,
|
|
pub(super) target_exec: String,
|
|
pub(super) sandbox_ready_at: Option<u64>,
|
|
pub(super) exec_established_at: Option<u64>,
|
|
pub(super) source_fingerprint_before: String,
|
|
pub(super) started_at: u64,
|
|
pub(super) output: Mutex<ProcessOutputState>,
|
|
pub(super) output_changed: Condvar,
|
|
pub(super) writer: Mutex<Option<Box<dyn std::io::Write + Send>>>,
|
|
pub(super) master: Mutex<Option<Box<dyn MasterPty + Send>>>,
|
|
#[cfg(windows)]
|
|
pub(super) job: Mutex<Option<WindowsProcessJob>>,
|
|
pub(super) control: std::sync::mpsc::Sender<ProcessControl>,
|
|
}
|
|
|
|
#[cfg(windows)]
|
|
#[derive(Debug)]
|
|
pub(crate) struct WindowsProcessJob(windows_sys::Win32::Foundation::HANDLE);
|
|
|
|
#[cfg(windows)]
|
|
unsafe impl Send for WindowsProcessJob {}
|
|
|
|
#[cfg(windows)]
|
|
unsafe impl Sync for WindowsProcessJob {}
|
|
|
|
#[cfg(windows)]
|
|
impl WindowsProcessJob {
|
|
pub(super) fn assign(child: &dyn Child) -> Result<Self, String> {
|
|
let process = child
|
|
.as_raw_handle()
|
|
.ok_or_else(|| "command.start Windows child 缺少 process handle".to_string())?
|
|
as windows_sys::Win32::Foundation::HANDLE;
|
|
Self::assign_handle(process)
|
|
}
|
|
|
|
pub(crate) fn assign_std(child: &std::process::Child) -> Result<Self, String> {
|
|
use std::os::windows::io::AsRawHandle;
|
|
Self::assign_handle(
|
|
AsRawHandle::as_raw_handle(child) as windows_sys::Win32::Foundation::HANDLE
|
|
)
|
|
}
|
|
|
|
pub(crate) fn assign_tokio(child: &tokio::process::Child) -> Result<Self, String> {
|
|
let handle = child
|
|
.raw_handle()
|
|
.ok_or_else(|| "子进程缺少 Windows process handle".to_string())?;
|
|
Self::assign_handle(handle as windows_sys::Win32::Foundation::HANDLE)
|
|
}
|
|
|
|
pub(crate) fn assign_tokio_named(
|
|
child: &tokio::process::Child,
|
|
name: &str,
|
|
) -> Result<Self, String> {
|
|
let name = Self::job_name(name)?;
|
|
let handle = child
|
|
.raw_handle()
|
|
.ok_or_else(|| "子进程缺少 Windows process handle".to_string())?;
|
|
Self::assign_handle_named(
|
|
handle as windows_sys::Win32::Foundation::HANDLE,
|
|
Some(&name),
|
|
)
|
|
}
|
|
|
|
/// 仅恢复刚以 CREATE_SUSPENDED 创建并已绑定本 Job 的唯一主线程。
|
|
pub(crate) fn resume_suspended_tokio(
|
|
&self,
|
|
child: &tokio::process::Child,
|
|
) -> Result<(), String> {
|
|
use windows_sys::Win32::Foundation::{CloseHandle, INVALID_HANDLE_VALUE};
|
|
use windows_sys::Win32::System::Diagnostics::ToolHelp::{
|
|
CreateToolhelp32Snapshot, Thread32First, Thread32Next, TH32CS_SNAPTHREAD, THREADENTRY32,
|
|
};
|
|
use windows_sys::Win32::System::JobObjects::IsProcessInJob;
|
|
use windows_sys::Win32::System::Threading::{
|
|
GetProcessIdOfThread, OpenThread, ResumeThread, THREAD_QUERY_LIMITED_INFORMATION,
|
|
THREAD_SUSPEND_RESUME,
|
|
};
|
|
let pid = child.id().ok_or("受控命令缺少进程身份")?;
|
|
let process = child.raw_handle().ok_or("受控命令缺少进程句柄")?
|
|
as windows_sys::Win32::Foundation::HANDLE;
|
|
let mut in_job = 0;
|
|
if unsafe { IsProcessInJob(process, self.0, &mut in_job) } == 0 || in_job == 0 {
|
|
return Err("拒绝恢复不属于本 Job 的进程".into());
|
|
}
|
|
let snapshot = unsafe { CreateToolhelp32Snapshot(TH32CS_SNAPTHREAD, 0) };
|
|
if snapshot == INVALID_HANDLE_VALUE {
|
|
return Err("读取受控命令主线程失败".into());
|
|
}
|
|
let mut entry = THREADENTRY32::default();
|
|
entry.dwSize = std::mem::size_of::<THREADENTRY32>() as u32;
|
|
let mut found = Vec::new();
|
|
let mut available = unsafe { Thread32First(snapshot, &mut entry) };
|
|
while available != 0 {
|
|
if entry.th32OwnerProcessID == pid {
|
|
found.push(entry.th32ThreadID);
|
|
}
|
|
available = unsafe { Thread32Next(snapshot, &mut entry) };
|
|
}
|
|
let enumeration_error = std::io::Error::last_os_error().raw_os_error();
|
|
unsafe { CloseHandle(snapshot) };
|
|
if enumeration_error != Some(windows_sys::Win32::Foundation::ERROR_NO_MORE_FILES as i32) {
|
|
return Err("受控命令线程枚举未完成,拒绝恢复".into());
|
|
}
|
|
if found.len() != 1 {
|
|
return Err("受控命令主线程身份不唯一,拒绝恢复".into());
|
|
}
|
|
let thread = unsafe {
|
|
OpenThread(
|
|
THREAD_SUSPEND_RESUME | THREAD_QUERY_LIMITED_INFORMATION,
|
|
0,
|
|
found[0],
|
|
)
|
|
};
|
|
if thread.is_null() {
|
|
return Err("打开受控命令主线程失败".into());
|
|
}
|
|
let belongs = unsafe { GetProcessIdOfThread(thread) } == pid;
|
|
let previous = if belongs {
|
|
unsafe { ResumeThread(thread) }
|
|
} else {
|
|
u32::MAX
|
|
};
|
|
unsafe { CloseHandle(thread) };
|
|
if previous != 1 {
|
|
return Err("受控命令恢复未确认,不能无门执行".into());
|
|
}
|
|
Ok(())
|
|
}
|
|
|
|
fn job_name(name: &str) -> Result<Vec<u16>, String> {
|
|
if !name.starts_with("Local\\AGC")
|
|
|| name.len() > 160
|
|
|| !name
|
|
.bytes()
|
|
.all(|b| b.is_ascii_alphanumeric() || matches!(b, b'\\' | b'-'))
|
|
{
|
|
return Err("Windows Job 名称无效".into());
|
|
}
|
|
Ok(name.encode_utf16().chain(Some(0)).collect())
|
|
}
|
|
|
|
pub(crate) fn named_is_empty(name: &str) -> Result<bool, String> {
|
|
use windows_sys::Win32::Foundation::{CloseHandle, ERROR_FILE_NOT_FOUND};
|
|
use windows_sys::Win32::System::JobObjects::OpenJobObjectW;
|
|
// windows-sys 0.61 在 SystemServices 定义此 SDK 常量;无需为常量扩大 feature。
|
|
const JOB_OBJECT_QUERY: u32 = 0x0004;
|
|
let name = Self::job_name(name)?;
|
|
let handle = unsafe { OpenJobObjectW(JOB_OBJECT_QUERY, 0, name.as_ptr()) };
|
|
if handle.is_null() {
|
|
let error = std::io::Error::last_os_error();
|
|
return if error.raw_os_error() == Some(ERROR_FILE_NOT_FOUND as i32) {
|
|
Ok(true)
|
|
} else {
|
|
Err("无法核对 Windows Job 所有权".into())
|
|
};
|
|
}
|
|
let result = Self::handle_is_empty(handle);
|
|
unsafe { CloseHandle(handle) };
|
|
result
|
|
}
|
|
|
|
pub(crate) fn is_empty(&self) -> Result<bool, String> {
|
|
Self::handle_is_empty(self.0)
|
|
}
|
|
|
|
fn handle_is_empty(handle: windows_sys::Win32::Foundation::HANDLE) -> Result<bool, String> {
|
|
use windows_sys::Win32::System::JobObjects::{
|
|
JobObjectBasicAccountingInformation, QueryInformationJobObject,
|
|
JOBOBJECT_BASIC_ACCOUNTING_INFORMATION,
|
|
};
|
|
let mut information = JOBOBJECT_BASIC_ACCOUNTING_INFORMATION::default();
|
|
let okay = unsafe {
|
|
QueryInformationJobObject(
|
|
handle,
|
|
JobObjectBasicAccountingInformation,
|
|
&mut information as *mut _ as *mut _,
|
|
std::mem::size_of_val(&information) as u32,
|
|
std::ptr::null_mut(),
|
|
)
|
|
};
|
|
if okay == 0 {
|
|
return Err("无法确认 Windows Job 子树已退出".into());
|
|
}
|
|
Ok(information.ActiveProcesses == 0)
|
|
}
|
|
|
|
fn assign_handle(process: windows_sys::Win32::Foundation::HANDLE) -> Result<Self, String> {
|
|
Self::assign_handle_named(process, None)
|
|
}
|
|
|
|
fn assign_handle_named(
|
|
process: windows_sys::Win32::Foundation::HANDLE,
|
|
name: Option<&[u16]>,
|
|
) -> Result<Self, String> {
|
|
use std::mem::size_of;
|
|
use windows_sys::Win32::Foundation::{CloseHandle, INVALID_HANDLE_VALUE};
|
|
use windows_sys::Win32::System::JobObjects::{
|
|
AssignProcessToJobObject, CreateJobObjectW, JobObjectExtendedLimitInformation,
|
|
SetInformationJobObject, JOBOBJECT_EXTENDED_LIMIT_INFORMATION,
|
|
JOB_OBJECT_LIMIT_KILL_ON_JOB_CLOSE,
|
|
};
|
|
|
|
if name.is_some() {
|
|
unsafe { windows_sys::Win32::Foundation::SetLastError(0) };
|
|
}
|
|
let handle = unsafe {
|
|
CreateJobObjectW(
|
|
std::ptr::null(),
|
|
name.map_or(std::ptr::null(), |name| name.as_ptr()),
|
|
)
|
|
};
|
|
if handle.is_null() || handle == INVALID_HANDLE_VALUE {
|
|
return Err(format!(
|
|
"创建 command.start Windows Job Object 失败:{}",
|
|
std::io::Error::last_os_error()
|
|
));
|
|
}
|
|
if name.is_some()
|
|
&& unsafe { windows_sys::Win32::Foundation::GetLastError() }
|
|
== windows_sys::Win32::Foundation::ERROR_ALREADY_EXISTS
|
|
{
|
|
unsafe { CloseHandle(handle) };
|
|
return Err("Windows Job 身份已存在,拒绝合并进程".into());
|
|
}
|
|
let mut information = JOBOBJECT_EXTENDED_LIMIT_INFORMATION::default();
|
|
information.BasicLimitInformation.LimitFlags = JOB_OBJECT_LIMIT_KILL_ON_JOB_CLOSE;
|
|
let configured = unsafe {
|
|
SetInformationJobObject(
|
|
handle,
|
|
JobObjectExtendedLimitInformation,
|
|
&information as *const _ as *const _,
|
|
size_of::<JOBOBJECT_EXTENDED_LIMIT_INFORMATION>() as u32,
|
|
)
|
|
};
|
|
let assigned = configured != 0 && unsafe { AssignProcessToJobObject(handle, process) } != 0;
|
|
if !assigned {
|
|
let error = std::io::Error::last_os_error();
|
|
unsafe {
|
|
CloseHandle(handle);
|
|
}
|
|
return Err(format!(
|
|
"配置 command.start Windows Job Object 失败:{error}"
|
|
));
|
|
}
|
|
Ok(Self(handle))
|
|
}
|
|
|
|
pub(crate) fn terminate(&self) -> Result<(), String> {
|
|
use windows_sys::Win32::System::JobObjects::TerminateJobObject;
|
|
if unsafe { TerminateJobObject(self.0, 1) } == 0 {
|
|
return Err(format!(
|
|
"终止 command.start Windows Job Object 失败:{}",
|
|
std::io::Error::last_os_error()
|
|
));
|
|
}
|
|
Ok(())
|
|
}
|
|
}
|
|
|
|
#[cfg(windows)]
|
|
impl Drop for WindowsProcessJob {
|
|
fn drop(&mut self) {
|
|
unsafe {
|
|
windows_sys::Win32::Foundation::CloseHandle(self.0);
|
|
}
|
|
}
|
|
}
|
|
|
|
impl std::fmt::Debug for LiveProcessSession {
|
|
fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
|
|
formatter
|
|
.debug_struct("LiveProcessSession")
|
|
.field("process_id", &self.process_id)
|
|
.field("agent_id", &self.identity.agent_id)
|
|
.field("run_id", &self.identity.run_id)
|
|
.finish_non_exhaustive()
|
|
}
|
|
}
|
|
|
|
#[derive(Default)]
|
|
pub(super) struct ProcessSessionRegistry {
|
|
pub(super) sessions: HashMap<String, Arc<LiveProcessSession>>,
|
|
}
|
|
|
|
#[cfg(target_os = "linux")]
|
|
#[derive(Clone, Debug)]
|
|
pub(super) struct PendingProcessLaunch {
|
|
pub(super) root: PathBuf,
|
|
pub(super) agent_id: String,
|
|
pub(super) process_group_leader: Option<i32>,
|
|
pub(super) shutdown_requested: bool,
|
|
}
|
|
|
|
#[cfg(target_os = "linux")]
|
|
#[derive(Default)]
|
|
pub(super) struct PendingProcessLaunchRegistry {
|
|
pub(super) launches: HashMap<String, PendingProcessLaunch>,
|
|
}
|
|
|
|
#[cfg(target_os = "linux")]
|
|
pub(super) struct PendingProcessLaunchGuard {
|
|
process_id: String,
|
|
}
|
|
|
|
#[cfg(target_os = "linux")]
|
|
impl Drop for PendingProcessLaunchGuard {
|
|
fn drop(&mut self) {
|
|
if let Ok(mut registry) = pending_process_launch_registry().lock() {
|
|
registry.launches.remove(&self.process_id);
|
|
}
|
|
}
|
|
}
|
|
|
|
static PROCESS_SESSION_REGISTRY: OnceLock<Mutex<ProcessSessionRegistry>> = OnceLock::new();
|
|
static PROCESS_SESSION_BOOT_ID: OnceLock<String> = OnceLock::new();
|
|
#[cfg(target_os = "linux")]
|
|
static PENDING_PROCESS_LAUNCH_REGISTRY: OnceLock<Mutex<PendingProcessLaunchRegistry>> =
|
|
OnceLock::new();
|
|
|
|
pub(super) fn process_session_registry() -> &'static Mutex<ProcessSessionRegistry> {
|
|
PROCESS_SESSION_REGISTRY.get_or_init(|| Mutex::new(ProcessSessionRegistry::default()))
|
|
}
|
|
|
|
#[cfg(target_os = "linux")]
|
|
pub(super) fn pending_process_launch_registry() -> &'static Mutex<PendingProcessLaunchRegistry> {
|
|
PENDING_PROCESS_LAUNCH_REGISTRY
|
|
.get_or_init(|| Mutex::new(PendingProcessLaunchRegistry::default()))
|
|
}
|
|
|
|
#[cfg(target_os = "linux")]
|
|
pub(super) fn reserve_pending_process_launch(
|
|
root: &Path,
|
|
agent_id: &str,
|
|
process_id: &str,
|
|
) -> Result<PendingProcessLaunchGuard, String> {
|
|
let mut registry = pending_process_launch_registry()
|
|
.lock()
|
|
.map_err(|_| "pending process launch registry 锁已损坏".to_string())?;
|
|
if registry.launches.contains_key(process_id) {
|
|
return Err("command.start pending launch 身份冲突".to_string());
|
|
}
|
|
registry.launches.insert(
|
|
process_id.to_string(),
|
|
PendingProcessLaunch {
|
|
root: root.to_path_buf(),
|
|
agent_id: agent_id.to_string(),
|
|
process_group_leader: None,
|
|
shutdown_requested: false,
|
|
},
|
|
);
|
|
Ok(PendingProcessLaunchGuard {
|
|
process_id: process_id.to_string(),
|
|
})
|
|
}
|
|
|
|
#[cfg(target_os = "linux")]
|
|
pub(super) fn activate_pending_process_launch(
|
|
process_id: &str,
|
|
process_group_leader: i32,
|
|
) -> Result<(), String> {
|
|
if process_group_leader <= 1 {
|
|
return Err("command.start wrapper 进程组身份无效".to_string());
|
|
}
|
|
let mut registry = pending_process_launch_registry()
|
|
.lock()
|
|
.map_err(|_| "pending process launch registry 锁已损坏".to_string())?;
|
|
let launch = registry
|
|
.launches
|
|
.get_mut(process_id)
|
|
.ok_or_else(|| "command.start pending launch reservation 缺失".to_string())?;
|
|
launch.process_group_leader = Some(process_group_leader);
|
|
if launch.shutdown_requested {
|
|
unsafe {
|
|
libc::kill(-process_group_leader, libc::SIGKILL);
|
|
}
|
|
return Err("Runner shutdown 已取消 pending process launch".to_string());
|
|
}
|
|
Ok(())
|
|
}
|
|
|
|
pub(crate) fn process_session_boot_id() -> &'static str {
|
|
PROCESS_SESSION_BOOT_ID
|
|
.get_or_init(|| {
|
|
let mut digest = Sha256::new();
|
|
digest.update(std::process::id().to_le_bytes());
|
|
digest.update(
|
|
SystemTime::now()
|
|
.duration_since(UNIX_EPOCH)
|
|
.unwrap_or_default()
|
|
.as_nanos()
|
|
.to_le_bytes(),
|
|
);
|
|
digest.update(unix_timestamp().to_le_bytes());
|
|
let value = format!("{:x}", digest.finalize());
|
|
format!("boot-{}", &value[..32])
|
|
})
|
|
.as_str()
|
|
}
|
|
|
|
pub(crate) fn initialize_process_session_boot_id(boot_id: &str) -> Result<(), String> {
|
|
let boot_id = boot_id.trim();
|
|
if boot_id.is_empty()
|
|
|| boot_id.chars().count() > 160
|
|
|| boot_id.chars().any(|character| character.is_control())
|
|
{
|
|
return Err("Agent Runner bootId 无效,无法初始化 process session owner".to_string());
|
|
}
|
|
match PROCESS_SESSION_BOOT_ID.set(boot_id.to_string()) {
|
|
Ok(()) => Ok(()),
|
|
Err(_) if PROCESS_SESSION_BOOT_ID.get().map(String::as_str) == Some(boot_id) => Ok(()),
|
|
Err(_) => Err("process session owner bootId 已被其他 Runner 初始化".to_string()),
|
|
}
|
|
}
|
|
|
|
#[cfg(target_os = "linux")]
|
|
pub(crate) fn is_process_session_child_mode(args: &[String]) -> bool {
|
|
args.first().map(String::as_str) == Some(PROCESS_SESSION_CHILD_MODE)
|
|
}
|
|
|
|
#[cfg(target_os = "linux")]
|
|
pub(crate) fn run_process_session_child(args: &[String]) -> Result<i32, String> {
|
|
if args != [PROCESS_SESSION_CHILD_MODE] {
|
|
return Err("process session child 参数无效".to_string());
|
|
}
|
|
let expected_parent = std::env::var(PROCESS_SESSION_OWNER_PID_ENV)
|
|
.map_err(|_| "process session child 缺少 owner pid".to_string())?
|
|
.parse::<libc::pid_t>()
|
|
.map_err(|_| "process session child owner pid 无效".to_string())?;
|
|
if expected_parent <= 1 {
|
|
return Err("process session child owner pid 无效".to_string());
|
|
}
|
|
if unsafe { libc::getppid() } != expected_parent {
|
|
return Err("process session owner 在 child containment 生效前已退出".to_string());
|
|
}
|
|
unsafe {
|
|
libc::signal(libc::SIGHUP, libc::SIG_IGN);
|
|
}
|
|
std::env::remove_var(PROCESS_SESSION_OWNER_PID_ENV);
|
|
run_process_session_bridge_child(expected_parent)
|
|
}
|