feat: 完成 OpenNexus 第三阶段核心功能与生产化基础 #45
@@ -369,3 +369,13 @@ Core 的独立数据目录目前不等于已授权 Vault。Python 旧笔记写
|
|||||||
- 原生参数/环境探针、IPv4/IPv6 环回探针已切换到此正式进程模块执行,相关行为回归通过。新增实际 Windows 测试验证不匹配 SID 拒绝、丢弃未恢复进程后句柄进入终止状态、入口不存在与损坏终止符失败不消耗 Profile。运行中的等待探针验证零时限返回未完成、超过 60 秒的等待请求拒绝、主动终止后五秒内观察到退出。
|
- 原生参数/环境探针、IPv4/IPv6 环回探针已切换到此正式进程模块执行,相关行为回归通过。新增实际 Windows 测试验证不匹配 SID 拒绝、丢弃未恢复进程后句柄进入终止状态、入口不存在与损坏终止符失败不消耗 Profile。运行中的等待探针验证零时限返回未完成、超过 60 秒的等待请求拒绝、主动终止后五秒内观察到退出。
|
||||||
- 47 项扩展回归通过,3 项 ignored 为由父测试驱动的 Job 辅助入口;全目标 Clippy -D warnings 通过。日志 `.build/extension-process-regression-tests.log`、`.build/extension-process-clippy.log`。
|
- 47 项扩展回归通过,3 项 ignored 为由父测试驱动的 Job 辅助入口;全目标 Clippy -D warnings 通过。日志 `.build/extension-process-regression-tests.log`、`.build/extension-process-clippy.log`。
|
||||||
- 该模块只收拢已验证的原生生命周期,不构成完整生产扩展启动器。许可和包句柄原子绑定、撤销停止、生产 broker、scratch 配额、工具总期限及崩溃后的配置清理仍未完成,整体生产化继续实施。
|
- 该模块只收拢已验证的原生生命周期,不构成完整生产扩展启动器。许可和包句柄原子绑定、撤销停止、生产 broker、scratch 配额、工具总期限及崩溃后的配置清理仍未完成,整体生产化继续实施。
|
||||||
|
|
||||||
|
|
||||||
|
## 增量:独立于调用方轮询的单工具期限守卫
|
||||||
|
|
||||||
|
- 新增 ToolDeadline,生产预算固定 60 秒,按单次工具调用启动,不限制常驻 MCP 服务总寿命。Running::start_tool_call 绑定该实例的 Job,创建独立计时线程;线程创建或句柄复制失败返回明确错误,调用方必须在派发工具请求前成功启动守卫。
|
||||||
|
- 期限使用单调时钟和条件变量。超时独立终止整个实例 Job,不依赖调用方继续轮询;完成只能在期限到达前撤销计时,完成操作不能在期限已过、计时线程尚未调度时绕过超时。终止失败、计时线程异常与普通超时分别返回稳定错误码。
|
||||||
|
- finish 正常完成后保留常驻进程;cancel 主动终止实例;守卫未 finish 就被释放也终止实例,覆盖调用 future 取消或异常退出,避免留下无期限工具调用。释放会唤醒并回收计时线程,正常完成无需等待满 60 秒。
|
||||||
|
- 真实 AppContainer 等待探针验证:正常完成守卫后进程仍活着;使用测试内部 100 ms 期限时,在没有 check/finish 轮询驱动的情况下进程退出,随后 finish 明确返回超时;显式取消和未完成释放均在 5 秒内观察到进程退出。生产公开 API 不接受延长期限参数。本轮没有执行完整 60 秒墙钟超限、计时线程故障注入或高并发调用压力测试,不能据此独立通过 C-04。
|
||||||
|
- 48 项扩展回归通过、3 项 Job 辅助入口 ignored(由父测试执行),全目标 Clippy -D warnings 通过。日志 `.build/extension-deadline-tests.log`、`.build/extension-deadline-clippy.log`。
|
||||||
|
- 守卫仍需接到实际 MCP/扩展调用派发与结束路径;正式许可/包句柄绑定、broker、scratch 配额及完整验收尚未完成,整体生产化目标继续进行。
|
||||||
|
|||||||
@@ -707,12 +707,51 @@ mod tests {
|
|||||||
.code,
|
.code,
|
||||||
"EXTENSION_PROCESS_WAIT_INVALID"
|
"EXTENSION_PROCESS_WAIT_INVALID"
|
||||||
);
|
);
|
||||||
running.terminate().unwrap();
|
let completed = running.start_tool_call().unwrap();
|
||||||
|
completed.finish().unwrap();
|
||||||
|
assert_eq!(running.wait(std::time::Duration::ZERO).unwrap(), None);
|
||||||
|
let deadline = running
|
||||||
|
.start_test_tool_call(std::time::Duration::from_millis(100))
|
||||||
|
.unwrap();
|
||||||
|
// No check() or finish() drives expiration; wait only on the OS process.
|
||||||
|
assert!(running
|
||||||
|
.wait(std::time::Duration::from_secs(5))
|
||||||
|
.unwrap()
|
||||||
|
.is_some());
|
||||||
|
assert_eq!(
|
||||||
|
deadline.finish().unwrap_err().code,
|
||||||
|
"EXTENSION_TOOL_DEADLINE_EXCEEDED"
|
||||||
|
);
|
||||||
assert!(running
|
assert!(running
|
||||||
.wait(std::time::Duration::from_secs(5))
|
.wait(std::time::Duration::from_secs(5))
|
||||||
.unwrap()
|
.unwrap()
|
||||||
.is_some());
|
.is_some());
|
||||||
drop(running);
|
drop(running);
|
||||||
|
for explicit_cancel in [false, true] {
|
||||||
|
let data = crate::extension_launch_data::LaunchData::new(
|
||||||
|
&executable,
|
||||||
|
&["wait".to_owned()],
|
||||||
|
&system,
|
||||||
|
&folder,
|
||||||
|
&folder.join("Temp"),
|
||||||
|
&std::collections::BTreeMap::new(),
|
||||||
|
)
|
||||||
|
.unwrap();
|
||||||
|
let suspended =
|
||||||
|
crate::extension_process::Suspended::create(&profile, &executable, data).unwrap();
|
||||||
|
let running = unsafe { suspended.resume().unwrap() };
|
||||||
|
let deadline = running.start_tool_call().unwrap();
|
||||||
|
assert_eq!(running.wait(std::time::Duration::ZERO).unwrap(), None);
|
||||||
|
if explicit_cancel {
|
||||||
|
deadline.cancel().unwrap();
|
||||||
|
} else {
|
||||||
|
drop(deadline);
|
||||||
|
}
|
||||||
|
assert!(running
|
||||||
|
.wait(std::time::Duration::from_secs(5))
|
||||||
|
.unwrap()
|
||||||
|
.is_some());
|
||||||
|
}
|
||||||
for ip in ["127.0.0.1:0", "[::1]:0"] {
|
for ip in ["127.0.0.1:0", "[::1]:0"] {
|
||||||
let tcp = TcpListener::bind(ip).unwrap();
|
let tcp = TcpListener::bind(ip).unwrap();
|
||||||
let udp = UdpSocket::bind(ip).unwrap();
|
let udp = UdpSocket::bind(ip).unwrap();
|
||||||
|
|||||||
@@ -0,0 +1,158 @@
|
|||||||
|
//! Per-tool-call deadline. A persistent server does not have a 60-second lifetime.
|
||||||
|
use crate::{
|
||||||
|
extension_job::Job,
|
||||||
|
workspace::{HostError, Result},
|
||||||
|
};
|
||||||
|
use std::{
|
||||||
|
sync::{
|
||||||
|
atomic::{AtomicU8, Ordering},
|
||||||
|
Arc, Condvar, Mutex,
|
||||||
|
},
|
||||||
|
thread::JoinHandle,
|
||||||
|
time::{Duration, Instant},
|
||||||
|
};
|
||||||
|
|
||||||
|
pub const TOOL_BUDGET: Duration = Duration::from_secs(60);
|
||||||
|
struct State {
|
||||||
|
cancelled: Mutex<bool>,
|
||||||
|
wake: Condvar,
|
||||||
|
outcome: AtomicU8,
|
||||||
|
}
|
||||||
|
/// Host-owned guard, created before dispatching a tool call. Finish it only when
|
||||||
|
/// the call completes. Expiration kills the entire instance group even if the
|
||||||
|
/// caller never polls. Dropping without finish aborts the instance; only explicit
|
||||||
|
/// successful completion cancels the timer while preserving the server.
|
||||||
|
pub struct ToolDeadline {
|
||||||
|
state: Arc<State>,
|
||||||
|
job: Arc<Job>,
|
||||||
|
expires: Instant,
|
||||||
|
worker: Option<JoinHandle<()>>,
|
||||||
|
}
|
||||||
|
impl ToolDeadline {
|
||||||
|
pub(crate) fn arm(job: &Job) -> Result<Self> {
|
||||||
|
Self::with_budget(job, TOOL_BUDGET)
|
||||||
|
}
|
||||||
|
fn with_budget(job: &Job, budget: Duration) -> Result<Self> {
|
||||||
|
if budget.is_zero() || budget > TOOL_BUDGET {
|
||||||
|
return Err(HostError::new("EXTENSION_TOOL_BUDGET_INVALID"));
|
||||||
|
}
|
||||||
|
let expires = Instant::now() + budget;
|
||||||
|
let job = Arc::new(job.clone_for_deadline()?);
|
||||||
|
let state = Arc::new(State {
|
||||||
|
cancelled: Mutex::new(false),
|
||||||
|
wake: Condvar::new(),
|
||||||
|
outcome: AtomicU8::new(0),
|
||||||
|
});
|
||||||
|
let thread_job = Arc::clone(&job);
|
||||||
|
let thread_state = Arc::clone(&state);
|
||||||
|
let worker = std::thread::Builder::new()
|
||||||
|
.name("extension-tool-deadline".into())
|
||||||
|
.spawn(move || {
|
||||||
|
let mut cancelled = thread_state
|
||||||
|
.cancelled
|
||||||
|
.lock()
|
||||||
|
.unwrap_or_else(|e| e.into_inner());
|
||||||
|
loop {
|
||||||
|
if *cancelled {
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
let remaining = expires.saturating_duration_since(Instant::now());
|
||||||
|
if remaining.is_zero() {
|
||||||
|
break;
|
||||||
|
}
|
||||||
|
cancelled = thread_state
|
||||||
|
.wake
|
||||||
|
.wait_timeout(cancelled, remaining)
|
||||||
|
.unwrap_or_else(|e| e.into_inner())
|
||||||
|
.0;
|
||||||
|
}
|
||||||
|
thread_state.outcome.store(1, Ordering::Release);
|
||||||
|
drop(cancelled);
|
||||||
|
if thread_job.terminate().is_err() {
|
||||||
|
thread_state.outcome.store(2, Ordering::Release);
|
||||||
|
}
|
||||||
|
})
|
||||||
|
.map_err(|_| HostError::new("EXTENSION_TOOL_WATCHDOG_UNAVAILABLE"))?;
|
||||||
|
Ok(Self {
|
||||||
|
state,
|
||||||
|
job,
|
||||||
|
expires,
|
||||||
|
worker: Some(worker),
|
||||||
|
})
|
||||||
|
}
|
||||||
|
pub fn check(&self) -> Result<()> {
|
||||||
|
match self.state.outcome.load(Ordering::Acquire) {
|
||||||
|
2 => Err(HostError::new("EXTENSION_RESOURCE_TERMINATE_FAILED")),
|
||||||
|
3 => Err(HostError::new("EXTENSION_TOOL_WATCHDOG_FAILED")),
|
||||||
|
1 => Err(HostError::new("EXTENSION_TOOL_DEADLINE_EXCEEDED")),
|
||||||
|
_ if Instant::now() >= self.expires => {
|
||||||
|
Err(HostError::new("EXTENSION_TOOL_DEADLINE_EXCEEDED"))
|
||||||
|
}
|
||||||
|
_ => Ok(()),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
/// Completion cannot cancel an already elapsed budget, even if the timer
|
||||||
|
/// thread has not yet been scheduled to observe expiration.
|
||||||
|
pub fn finish(mut self) -> Result<()> {
|
||||||
|
self.stop();
|
||||||
|
match self.state.outcome.load(Ordering::Acquire) {
|
||||||
|
0 => Ok(()),
|
||||||
|
2 => Err(HostError::new("EXTENSION_RESOURCE_TERMINATE_FAILED")),
|
||||||
|
3 => Err(HostError::new("EXTENSION_TOOL_WATCHDOG_FAILED")),
|
||||||
|
_ => Err(HostError::new("EXTENSION_TOOL_DEADLINE_EXCEEDED")),
|
||||||
|
}
|
||||||
|
}
|
||||||
|
/// Explicit abandonment terminates the instance and reports kill failures.
|
||||||
|
pub fn cancel(mut self) -> Result<()> {
|
||||||
|
let result = self.job.terminate();
|
||||||
|
self.stop();
|
||||||
|
result
|
||||||
|
}
|
||||||
|
fn stop(&mut self) {
|
||||||
|
if self.worker.is_none() {
|
||||||
|
return;
|
||||||
|
}
|
||||||
|
let mut cancelled = self
|
||||||
|
.state
|
||||||
|
.cancelled
|
||||||
|
.lock()
|
||||||
|
.unwrap_or_else(|e| e.into_inner());
|
||||||
|
if Instant::now() < self.expires {
|
||||||
|
*cancelled = true;
|
||||||
|
}
|
||||||
|
self.state.wake.notify_all();
|
||||||
|
drop(cancelled);
|
||||||
|
if self.worker.take().unwrap().join().is_err() {
|
||||||
|
let _ = self.job.terminate();
|
||||||
|
self.state.outcome.store(3, Ordering::Release);
|
||||||
|
}
|
||||||
|
}
|
||||||
|
#[cfg(test)]
|
||||||
|
pub(crate) fn arm_test(job: &Job, budget: Duration) -> Result<Self> {
|
||||||
|
Self::with_budget(job, budget)
|
||||||
|
}
|
||||||
|
}
|
||||||
|
impl Drop for ToolDeadline {
|
||||||
|
fn drop(&mut self) {
|
||||||
|
if self.worker.is_some() {
|
||||||
|
let _ = self.job.terminate();
|
||||||
|
self.stop();
|
||||||
|
}
|
||||||
|
}
|
||||||
|
}
|
||||||
|
|
||||||
|
#[cfg(test)]
|
||||||
|
mod tests {
|
||||||
|
use super::*;
|
||||||
|
#[test]
|
||||||
|
fn rejects_invalid_budgets_and_finishes_without_waiting_for_full_budget() {
|
||||||
|
let job = Job::new().unwrap();
|
||||||
|
assert!(ToolDeadline::with_budget(&job, Duration::ZERO).is_err());
|
||||||
|
assert!(ToolDeadline::with_budget(&job, Duration::from_secs(61)).is_err());
|
||||||
|
let deadline = ToolDeadline::arm(&job).unwrap();
|
||||||
|
assert!(deadline.check().is_ok());
|
||||||
|
let start = Instant::now();
|
||||||
|
deadline.finish().unwrap();
|
||||||
|
assert!(start.elapsed() < Duration::from_secs(5));
|
||||||
|
}
|
||||||
|
}
|
||||||
@@ -13,6 +13,14 @@ pub struct Job {
|
|||||||
handle: OwnedHandle,
|
handle: OwnedHandle,
|
||||||
}
|
}
|
||||||
impl Job {
|
impl Job {
|
||||||
|
pub(crate) fn clone_for_deadline(&self) -> Result<Self> {
|
||||||
|
Ok(Self {
|
||||||
|
handle: self
|
||||||
|
.handle
|
||||||
|
.try_clone()
|
||||||
|
.map_err(|_| HostError::new("EXTENSION_RESOURCE_UNAVAILABLE"))?,
|
||||||
|
})
|
||||||
|
}
|
||||||
pub fn new() -> Result<Self> {
|
pub fn new() -> Result<Self> {
|
||||||
Self::with_process_limit(16)
|
Self::with_process_limit(16)
|
||||||
}
|
}
|
||||||
|
|||||||
@@ -175,6 +175,18 @@ impl<'a> Suspended<'a> {
|
|||||||
}
|
}
|
||||||
}
|
}
|
||||||
impl Running<'_> {
|
impl Running<'_> {
|
||||||
|
/// Arm before dispatching a tool request; finish after receiving its result.
|
||||||
|
/// Failure to arm must prevent dispatch. This does not time server lifetime.
|
||||||
|
pub fn start_tool_call(&self) -> Result<crate::extension_deadline::ToolDeadline> {
|
||||||
|
crate::extension_deadline::ToolDeadline::arm(&self.0.job)
|
||||||
|
}
|
||||||
|
#[cfg(test)]
|
||||||
|
pub(crate) fn start_test_tool_call(
|
||||||
|
&self,
|
||||||
|
budget: Duration,
|
||||||
|
) -> Result<crate::extension_deadline::ToolDeadline> {
|
||||||
|
crate::extension_deadline::ToolDeadline::arm_test(&self.0.job, budget)
|
||||||
|
}
|
||||||
/// A bounded observation only. The runtime must enforce the tool deadline.
|
/// A bounded observation only. The runtime must enforce the tool deadline.
|
||||||
pub fn wait(&self, timeout: Duration) -> Result<Option<u32>> {
|
pub fn wait(&self, timeout: Duration) -> Result<Option<u32>> {
|
||||||
let milliseconds = u32::try_from(timeout.as_millis())
|
let milliseconds = u32::try_from(timeout.as_millis())
|
||||||
|
|||||||
@@ -63,3 +63,6 @@ pub mod extension_launch_data;
|
|||||||
|
|
||||||
#[cfg(windows)]
|
#[cfg(windows)]
|
||||||
pub mod extension_process;
|
pub mod extension_process;
|
||||||
|
|
||||||
|
#[cfg(windows)]
|
||||||
|
pub mod extension_deadline;
|
||||||
|
|||||||
Reference in New Issue
Block a user