| 1 | //! One admitted process operation. The Builtin orchestrates; Rust owns all launch and caller facts. |
| 2 | use std::collections::HashMap; |
| 3 | use std::sync::{ |
| 4 | Arc, Mutex, |
| 5 | atomic::{AtomicU64, Ordering}, |
| 6 | }; |
| 7 | use std::time::{Duration, Instant}; |
| 8 | |
| 9 | use serde_json::{Value, json}; |
| 10 | use tokio::sync::{OwnedSemaphorePermit, Semaphore}; |
| 11 | use tokio_util::sync::CancellationToken; |
| 12 | |
| 13 | use super::protocol::{ |
| 14 | CoreRequest, ExecutionRedeemParams, HarnessRunParams, OwnerRef, RpcErrorWire, error_code, |
| 15 | }; |
| 16 | use super::registry::OwnerState; |
| 17 | use super::supervisor::HostRequestContext; |
| 18 | use super::ticket::{Grant, Presented, Ticket, TicketKind}; |
| 19 | use super::tier::HostTier; |
| 20 | use super::{ExtensionHostManager, ManagerShared, activation}; |
| 21 | use crate::plugins::PluginRegistry; |
| 22 | use crate::plugins::activation::PluginActivationCapability; |
| 23 | use crate::plugins::types::PluginAuthority; |
| 24 | use crate::tools::spec::{ToolContext, ToolError, ToolResult}; |
| 25 | |
| 26 | const DEADLINE: Duration = Duration::from_secs(120); |
| 27 | const MAX_JOBS: usize = 32; |
| 28 | const MAX_INPUT: usize = 1024 * 1024; |
| 29 | const MAX_REVIEW: usize = 16 * 1024 * 1024; |
| 30 | const MAX_STDOUT: usize = 1024 * 1024; |
| 31 | const MAX_STDERR: usize = 64 * 1024; |
| 32 | |
| 33 | #[derive(Clone)] |
| 34 | struct Caller { |
| 35 | workspace: std::path::PathBuf, |
| 36 | plugins: Option<Arc<PluginRegistry>>, |
| 37 | session_id: Option<String>, |
| 38 | agent_id: Option<String>, |
| 39 | origin_turn_id: Option<String>, |
| 40 | origin_call_id: Option<String>, |
| 41 | authorities: Vec<PluginAuthority>, |
| 42 | native_owners: Vec<OwnerRef>, |
| 43 | native_host_generation: Option<u64>, |
| 44 | requires_native_policy: bool, |
| 45 | } |
| 46 | impl Caller { |
| 47 | fn capture(context: &ToolContext, shared: &ManagerShared) -> Result<Self, String> { |
| 48 | Self::from_hook(crate::hooks::HookCaller::from_tool(context), shared, None) |
| 49 | } |
| 50 | fn from_hook( |
| 51 | caller: crate::hooks::HookCaller, |
| 52 | shared: &ManagerShared, |
| 53 | only_plugin: Option<&str>, |
| 54 | ) -> Result<Self, String> { |
| 55 | let plugins = caller.plugins.clone(); |
| 56 | let mut authorities = Vec::new(); |
| 57 | if let Some(plugins) = plugins.as_ref() { |
| 58 | for entry in plugins |
| 59 | .selected_native_entries() |
| 60 | .iter() |
| 61 | .filter(|entry| only_plugin.is_none_or(|id| entry.plugin_id == id)) |
| 62 | { |
| 63 | let authority = plugins |
| 64 | .authority_for(&entry.plugin_id) |
| 65 | .ok_or("selected Native source has no authority")?; |
| 66 | if authority.content_hash != entry.content_hash { |
| 67 | return Err("selected Native source changed".into()); |
| 68 | } |
| 69 | if !authorities |
| 70 | .iter() |
| 71 | .any(|old: &PluginAuthority| old.plugin_id == authority.plugin_id) |
| 72 | { |
| 73 | authorities.push(authority); |
| 74 | } |
| 75 | } |
| 76 | } |
| 77 | let native_host_generation = if authorities.is_empty() { |
| 78 | None |
| 79 | } else { |
| 80 | Some( |
| 81 | shared |
| 82 | .ready_host(HostTier::Plugin) |
| 83 | .map_err(|status| status.to_string())? |
| 84 | .generation, |
| 85 | ) |
| 86 | }; |
| 87 | let native_owners = { |
| 88 | let registry = shared.registry.lock().expect("registry lock"); |
| 89 | authorities |
| 90 | .iter() |
| 91 | .map(|authority| { |
| 92 | registry |
| 93 | .owner(authority.plugin_id.as_str()) |
| 94 | .filter(|owner| { |
| 95 | owner.state == OwnerState::Active |
| 96 | && owner.content_hash == authority.content_hash |
| 97 | }) |
| 98 | .map(|owner| owner.owner.clone()) |
| 99 | .ok_or_else(|| "Native caller owner is no longer live".to_string()) |
| 100 | }) |
| 101 | .collect::<Result<Vec<_>, _>>()? |
| 102 | }; |
| 103 | Ok(Self { |
| 104 | workspace: caller.workspace, |
| 105 | plugins, |
| 106 | session_id: caller.session_id.filter(|id| !id.is_empty()), |
| 107 | agent_id: caller.agent_id, |
| 108 | origin_turn_id: caller.origin_turn_id, |
| 109 | origin_call_id: caller.origin_call_id, |
| 110 | authorities, |
| 111 | native_owners, |
| 112 | native_host_generation, |
| 113 | requires_native_policy: true, |
| 114 | }) |
| 115 | } |
| 116 | fn check(&self, shared: &ManagerShared) -> Result<(), String> { |
| 117 | if (self.requires_native_policy || !self.authorities.is_empty()) |
| 118 | && !activation::extension_host_policy_enabled() |
| 119 | { |
| 120 | return Err("extension host is disabled".into()); |
| 121 | } |
| 122 | if let Some(generation) = self.native_host_generation { |
| 123 | if shared |
| 124 | .ready_host(HostTier::Plugin) |
| 125 | .map_err(|status| status.to_string())? |
| 126 | .generation |
| 127 | != generation |
| 128 | { |
| 129 | return Err("Native caller host restarted".into()); |
| 130 | } |
| 131 | let registry = shared.registry.lock().expect("registry lock"); |
| 132 | for owner in &self.native_owners { |
| 133 | let active = registry |
| 134 | .owner(&owner.plugin_id) |
| 135 | .filter(|entry| entry.owner == *owner && entry.state == OwnerState::Active) |
| 136 | .ok_or("Native caller owner changed")?; |
| 137 | for entry in self |
| 138 | .plugins |
| 139 | .as_ref() |
| 140 | .expect("Native view") |
| 141 | .selected_native_entries() |
| 142 | .iter() |
| 143 | .filter(|entry| entry.plugin_id == owner.plugin_id) |
| 144 | { |
| 145 | if active.scopes.get(&entry.entry) != Some(&OwnerState::Active) { |
| 146 | return Err("Native caller entry was withdrawn".into()); |
| 147 | } |
| 148 | } |
| 149 | } |
| 150 | } |
| 151 | if let Some(plugins) = self.plugins.as_ref() { |
| 152 | if plugins.workspace() != self.workspace { |
| 153 | return Err("execution workspace does not match its caller".into()); |
| 154 | } |
| 155 | if let Some(revision) = plugins.caller_selection() { |
| 156 | let (current, agent) = shared |
| 157 | .plugins_for_selection(revision, self.session_id.as_deref()) |
| 158 | .ok_or("execution caller was detached or changed")?; |
| 159 | if agent != self.agent_id || !Arc::ptr_eq(plugins, ¤t) { |
| 160 | return Err("execution caller identity changed".into()); |
| 161 | } |
| 162 | } else if !plugins.selected_native_entries().is_empty() { |
| 163 | return Err("selected Native execution has no caller receipt".into()); |
| 164 | } |
| 165 | } |
| 166 | Ok(()) |
| 167 | } |
| 168 | fn target(&self, id: &str) -> Value { |
| 169 | json!({"execution_id":id,"workspace":self.workspace,"session_id":self.session_id,"agent_id":self.agent_id,"origin_turn_id":self.origin_turn_id,"origin_call_id":self.origin_call_id,"selection":self.plugins.as_ref().and_then(|plugins| plugins.caller_selection())}) |
| 170 | } |
| 171 | } |
| 172 | |
| 173 | /// An exhaustive pure-transform inventory, never an arbitrary module/function selector. |
| 174 | #[derive(Clone, Copy)] |
| 175 | pub(crate) enum StockOperation { |
| 176 | WebFilters, |
| 177 | WebRequest, |
| 178 | WebProvider, |
| 179 | WebEntries, |
| 180 | WebFinalize, |
| 181 | WebExtract, |
| 182 | WebImages, |
| 183 | FinanceQuote, |
| 184 | FinanceChart, |
| 185 | ValidateData, |
| 186 | SpeechOptions, |
| 187 | SpeechFormat, |
| 188 | ReviewSourcePrompt, |
| 189 | ReviewPassPrompt, |
| 190 | ReviewInteractivePr, |
| 191 | ReviewReport, |
| 192 | } |
| 193 | impl StockOperation { |
| 194 | fn as_str(self) -> &'static str { |
| 195 | match self { |
| 196 | Self::WebFilters => "web_filters", |
| 197 | Self::WebRequest => "web_request", |
| 198 | Self::WebProvider => "web_provider", |
| 199 | Self::WebEntries => "web_entries", |
| 200 | Self::WebFinalize => "web_finalize", |
| 201 | Self::WebExtract => "web_extract", |
| 202 | Self::WebImages => "web_images", |
| 203 | Self::FinanceQuote => "finance_quote", |
| 204 | Self::FinanceChart => "finance_chart", |
| 205 | Self::ValidateData => "validate_data", |
| 206 | Self::SpeechOptions => "speech_options", |
| 207 | Self::SpeechFormat => "speech_format", |
| 208 | Self::ReviewSourcePrompt => "review_source_prompt", |
| 209 | Self::ReviewPassPrompt => "review_pass_prompt", |
| 210 | Self::ReviewInteractivePr => "review_interactive_pr", |
| 211 | Self::ReviewReport => "review_report", |
| 212 | } |
| 213 | } |
| 214 | fn is_review(self) -> bool { |
| 215 | matches!( |
| 216 | self, |
| 217 | Self::ReviewSourcePrompt |
| 218 | | Self::ReviewPassPrompt |
| 219 | | Self::ReviewInteractivePr |
| 220 | | Self::ReviewReport |
| 221 | ) |
| 222 | } |
| 223 | fn slots(self) -> u32 { |
| 224 | if self.is_review() { 16 } else { 1 } |
| 225 | } |
| 226 | } |
| 227 | |
| 228 | fn review_envelope_fits(value: &Value) -> Result<bool, ToolError> { |
| 229 | super::protocol::encode_frame(&json!({"jsonrpc":"2.0","id":u64::MAX,"result":value})) |
| 230 | .map(|frame| frame.len() <= MAX_REVIEW) |
| 231 | .map_err(|_| { |
| 232 | ToolError::invalid_input("Serialized review envelope exceeds the existing frame limit") |
| 233 | }) |
| 234 | } |
| 235 | fn check_stock_snapshot(operation: StockOperation, input: &Value) -> Result<(), ToolError> { |
| 236 | if operation.is_review() { |
| 237 | let projection = |
| 238 | json!({"kind":"stock_adapter","operation":operation.as_str(),"input":input}); |
| 239 | if !review_envelope_fits(&projection)? { |
| 240 | return Err(ToolError::invalid_input( |
| 241 | "Serialized review snapshot/envelope exceeds 16 MiB; no source was truncated", |
| 242 | )); |
| 243 | } |
| 244 | } else if serde_json::to_vec(input) |
| 245 | .map_err(|_| ToolError::invalid_input("adapter snapshot is not JSON"))? |
| 246 | .len() |
| 247 | > MAX_INPUT |
| 248 | { |
| 249 | return Err(ToolError::invalid_input("adapter snapshot exceeds 1 MiB")); |
| 250 | } |
| 251 | Ok(()) |
| 252 | } |
| 253 | |
| 254 | type HookRun = Box<dyn FnOnce(CancellationToken) -> crate::hooks::HookResult + Send>; |
| 255 | struct HookLaunch { |
| 256 | run: HookRun, |
| 257 | hook: crate::hooks::Hook, |
| 258 | } |
| 259 | struct PdfLaunch { |
| 260 | input: crate::tools::pdf::CapturedPdf, |
| 261 | output: Arc<Mutex<Option<crate::tools::pdf::PdfProcessOutcome>>>, |
| 262 | } |
| 263 | enum OcrLaunch { |
| 264 | Native { |
| 265 | input: crate::tools::image_ocr::CapturedOcr, |
| 266 | output: Arc<Mutex<Option<crate::tools::image_ocr::OcrOutcome>>>, |
| 267 | }, |
| 268 | Tesseract { |
| 269 | step: crate::tools::image_ocr::OcrNativeStep, |
| 270 | output: Arc<Mutex<Option<crate::tools::image_ocr::OcrOutcome>>>, |
| 271 | }, |
| 272 | } |
| 273 | struct Continuation { |
| 274 | ticket: Ticket, |
| 275 | target: Value, |
| 276 | } |
| 277 | struct GithubLaunch { |
| 278 | request: crate::tools::github::host::Request, |
| 279 | context: Option<ToolContext>, |
| 280 | output: Arc<Mutex<crate::tools::github::host::RunState>>, |
| 281 | } |
| 282 | struct Launch { |
| 283 | github: Option<GithubLaunch>, |
| 284 | ocr: Option<OcrLaunch>, |
| 285 | pdf: Option<PdfLaunch>, |
| 286 | stock: Option<(StockOperation, Value)>, |
| 287 | hook: Option<HookLaunch>, |
| 288 | command: Option<tokio::process::Command>, |
| 289 | input: Vec<u8>, |
| 290 | } |
| 291 | struct Job { |
| 292 | id: String, |
| 293 | owner: OwnerRef, |
| 294 | generation: u64, |
| 295 | caller: Caller, |
| 296 | target: Value, |
| 297 | expires: Instant, |
| 298 | cancel: CancellationToken, |
| 299 | launch: Mutex<Option<Launch>>, |
| 300 | ticket: Ticket, |
| 301 | continuation: Mutex<Option<Continuation>>, |
| 302 | hook_result: Mutex<Option<crate::hooks::HookResult>>, |
| 303 | // The worker owns this Arc until process cleanup and blocking receipt checks finish. |
| 304 | _permit: OwnedSemaphorePermit, |
| 305 | } |
| 306 | impl Job { |
| 307 | fn check(&self, shared: &ManagerShared) -> Result<(), String> { |
| 308 | if self.cancel.is_cancelled() || Instant::now() >= self.expires { |
| 309 | return Err("execution cancelled or expired".into()); |
| 310 | } |
| 311 | if shared.builtin.host_generation.load(Ordering::SeqCst) != self.generation { |
| 312 | return Err("execution host restarted".into()); |
| 313 | } |
| 314 | shared.live_owner_authority(HostTier::Builtin, |registry| { |
| 315 | registry |
| 316 | .owner("host:harness") |
| 317 | .filter(|entry| entry.owner == self.owner && entry.state == OwnerState::Active) |
| 318 | .map(|entry| entry.owner.clone()) |
| 319 | .ok_or_else(|| "execution builtin owner is no longer live".to_string()) |
| 320 | })?; |
| 321 | self.caller.check(shared) |
| 322 | } |
| 323 | async fn checked(self: &Arc<Self>, shared: &Arc<ManagerShared>) -> Result<(), String> { |
| 324 | self.check(shared)?; |
| 325 | let job = Arc::clone(self); |
| 326 | let policy = activation::extension_host_policy_enabled(); |
| 327 | #[cfg(test)] |
| 328 | let env_scope = crate::test_support::env_scope_ticket(); |
| 329 | shared |
| 330 | .engine_handle()? |
| 331 | .spawn_blocking(move || { |
| 332 | let _policy = activation::PolicyScope::propagate(policy); |
| 333 | #[cfg(test)] |
| 334 | let _env_scope = crate::test_support::join_env_scope(env_scope); |
| 335 | for authority in &job.caller.authorities { |
| 336 | crate::plugins::registry::verify_plugin_component_authority( |
| 337 | authority, |
| 338 | PluginActivationCapability::Native, |
| 339 | )?; |
| 340 | } |
| 341 | Ok::<(), String>(()) |
| 342 | }) |
| 343 | .await |
| 344 | .map_err(|_| "execution receipt check failed".to_string())??; |
| 345 | self.check(shared) |
| 346 | } |
| 347 | } |
| 348 | |
| 349 | pub(super) struct Broker { |
| 350 | jobs: Mutex<HashMap<String, Arc<Job>>>, |
| 351 | slots: Arc<Semaphore>, |
| 352 | } |
| 353 | impl Default for Broker { |
| 354 | fn default() -> Self { |
| 355 | Self { |
| 356 | jobs: Mutex::new(HashMap::new()), |
| 357 | slots: Arc::new(Semaphore::new(MAX_JOBS)), |
| 358 | } |
| 359 | } |
| 360 | } |
| 361 | impl Drop for Broker { |
| 362 | fn drop(&mut self) { |
| 363 | for job in self.jobs.get_mut().expect("execution lock").values() { |
| 364 | job.cancel.cancel(); |
| 365 | } |
| 366 | } |
| 367 | } |
| 368 | impl Broker { |
| 369 | fn admit(&self, units: u32) -> Result<OwnedSemaphorePermit, ToolError> { |
| 370 | Arc::clone(&self.slots) |
| 371 | .try_acquire_many_owned(units) |
| 372 | .map_err(|_| ToolError::not_available("execution concurrency limit reached")) |
| 373 | } |
| 374 | pub(super) fn revoke_owner(&self, owner: &str) { |
| 375 | for job in self.jobs.lock().expect("execution lock").values() { |
| 376 | if job.owner.plugin_id == owner |
| 377 | || job |
| 378 | .caller |
| 379 | .authorities |
| 380 | .iter() |
| 381 | .any(|authority| authority.plugin_id.as_str() == owner) |
| 382 | { |
| 383 | job.cancel.cancel(); |
| 384 | } |
| 385 | } |
| 386 | } |
| 387 | pub(super) fn revoke_host(&self, tier: HostTier, generation: u64) { |
| 388 | for job in self.jobs.lock().expect("execution lock").values() { |
| 389 | if (tier == HostTier::Builtin && job.generation == generation) |
| 390 | || (tier == HostTier::Plugin |
| 391 | && job.caller.native_host_generation == Some(generation)) |
| 392 | { |
| 393 | job.cancel.cancel(); |
| 394 | } |
| 395 | } |
| 396 | } |
| 397 | pub(super) fn revoke_scope(&self, plugin_id: &str, scope: &super::protocol::EntryRef) { |
| 398 | for job in self.jobs.lock().expect("execution lock").values() { |
| 399 | if job.caller.plugins.as_ref().is_some_and(|plugins| { |
| 400 | plugins |
| 401 | .selected_native_entries() |
| 402 | .iter() |
| 403 | .any(|entry| entry.plugin_id == plugin_id && entry.entry == *scope) |
| 404 | }) { |
| 405 | job.cancel.cancel(); |
| 406 | } |
| 407 | } |
| 408 | } |
| 409 | pub(super) fn revoke_attachment(&self, id: u64) { |
| 410 | for job in self.jobs.lock().expect("execution lock").values() { |
| 411 | if job |
| 412 | .caller |
| 413 | .plugins |
| 414 | .as_ref() |
| 415 | .and_then(|plugins| plugins.caller_selection()) |
| 416 | .is_some_and(|selected| selected.attachment_id == id) |
| 417 | { |
| 418 | job.cancel.cancel(); |
| 419 | } |
| 420 | } |
| 421 | } |
| 422 | pub(super) async fn serve( |
| 423 | &self, |
| 424 | shared: &Arc<ManagerShared>, |
| 425 | generation: u64, |
| 426 | params: ExecutionRedeemParams, |
| 427 | cx: HostRequestContext, |
| 428 | ) -> Result<Value, RpcErrorWire> { |
| 429 | let refused = |message: String| RpcErrorWire { |
| 430 | code: error_code::REFUSED, |
| 431 | message, |
| 432 | data: None, |
| 433 | }; |
| 434 | if params.owner.plugin_id != "host:harness" || cx.cancel.is_cancelled() { |
| 435 | return Err(refused("execution is not admitted".into())); |
| 436 | } |
| 437 | let job = self |
| 438 | .jobs |
| 439 | .lock() |
| 440 | .expect("execution lock") |
| 441 | .get(¶ms.execution_id) |
| 442 | .cloned() |
| 443 | .ok_or_else(|| refused("execution is no longer admitted".into()))?; |
| 444 | if generation != job.generation || params.owner != job.owner { |
| 445 | return Err(refused("execution owner or generation changed".into())); |
| 446 | } |
| 447 | job.check(shared).map_err(refused)?; |
| 448 | let target = job |
| 449 | .continuation |
| 450 | .lock() |
| 451 | .expect("execution continuation lock") |
| 452 | .as_ref() |
| 453 | .map(|grant| grant.target.clone()) |
| 454 | .unwrap_or_else(|| job.target.clone()); |
| 455 | if let Err(error) = shared.core_calls.tickets.redeem(&Presented { |
| 456 | ticket: ¶ms.ticket, |
| 457 | kind: TicketKind::Execution, |
| 458 | tier: HostTier::Builtin, |
| 459 | host_generation: generation, |
| 460 | owner: ¶ms.owner, |
| 461 | method: "exec/redeem", |
| 462 | target: Some(&target), |
| 463 | }) { |
| 464 | if error.violation { |
| 465 | cx.violation("too many invalid execution tickets".to_string()); |
| 466 | } |
| 467 | return Err(RpcErrorWire { |
| 468 | code: error_code::REFUSED, |
| 469 | message: error.reason.describe().to_string(), |
| 470 | data: None, |
| 471 | }); |
| 472 | } |
| 473 | let mut guard = CancelOnDrop(job.cancel.clone()); |
| 474 | let shared_worker = Arc::clone(shared); |
| 475 | let job_worker = Arc::clone(&job); |
| 476 | let work = dispatch_job( |
| 477 | shared.engine_handle().map_err(refused)?, |
| 478 | job_worker, |
| 479 | move |job| run_job(shared_worker, job), |
| 480 | ); |
| 481 | let result = tokio::select! { |
| 482 | result = work => result.map_err(|_| refused("execution worker failed".into()))?, |
| 483 | _ = cx.cancel.cancelled() => { job.cancel.cancel(); return Err(refused("execution request cancelled".into())); } |
| 484 | }.map_err(refused)?; |
| 485 | job.check(shared).map_err(refused)?; |
| 486 | guard.0 = CancellationToken::new(); |
| 487 | Ok(result) |
| 488 | } |
| 489 | } |
| 490 | struct CancelOnDrop(CancellationToken); |
| 491 | impl Drop for CancelOnDrop { |
| 492 | fn drop(&mut self) { |
| 493 | self.0.cancel(); |
| 494 | } |
| 495 | } |
| 496 | struct Demand<'a>(&'a AtomicU64); |
| 497 | impl Drop for Demand<'_> { |
| 498 | fn drop(&mut self) { |
| 499 | self.0.fetch_sub(1, Ordering::SeqCst); |
| 500 | } |
| 501 | } |
| 502 | struct Invocation { |
| 503 | shared: Arc<ManagerShared>, |
| 504 | job: Arc<Job>, |
| 505 | } |
| 506 | impl Drop for Invocation { |
| 507 | fn drop(&mut self) { |
| 508 | self.job.cancel.cancel(); |
| 509 | self.shared.core_calls.tickets.revoke(&self.job.ticket); |
| 510 | if let Some(grant) = self |
| 511 | .job |
| 512 | .continuation |
| 513 | .lock() |
| 514 | .expect("execution continuation lock") |
| 515 | .as_ref() |
| 516 | { |
| 517 | self.shared.core_calls.tickets.revoke(&grant.ticket); |
| 518 | } |
| 519 | self.shared |
| 520 | .execution_broker |
| 521 | .jobs |
| 522 | .lock() |
| 523 | .expect("execution lock") |
| 524 | .remove(&self.job.id); |
| 525 | } |
| 526 | } |
| 527 | |
| 528 | // A dropped request awaits no result, but the Engine worker retains its job and quota |
| 529 | // until its admitted process/receipt cleanup completes. No transient runtime. |
| 530 | fn dispatch_job<F, R>( |
| 531 | handle: tokio::runtime::Handle, |
| 532 | job: Arc<Job>, |
| 533 | run: F, |
| 534 | ) -> tokio::task::JoinHandle<Result<Value, String>> |
| 535 | where |
| 536 | F: FnOnce(Arc<Job>) -> R + Send + 'static, |
| 537 | R: std::future::Future<Output = Result<Value, String>> + Send + 'static, |
| 538 | { |
| 539 | handle.spawn(async move { |
| 540 | let result = run(Arc::clone(&job)).await; |
| 541 | drop(job); |
| 542 | result |
| 543 | }) |
| 544 | } |
| 545 | |
| 546 | async fn run_job(shared: Arc<ManagerShared>, job: Arc<Job>) -> Result<Value, String> { |
| 547 | job.checked(&shared).await?; |
| 548 | let mut launch = job |
| 549 | .launch |
| 550 | .lock() |
| 551 | .expect("execution launch lock") |
| 552 | .take() |
| 553 | .ok_or("execution was already redeemed")?; |
| 554 | if let Some(github) = launch.github.take() { |
| 555 | let shared_check = Arc::clone(&shared); |
| 556 | let job_check = Arc::clone(&job); |
| 557 | let result = crate::tools::github::host::run( |
| 558 | github.request, |
| 559 | github.context, |
| 560 | job.cancel.clone(), |
| 561 | move || { |
| 562 | let shared = Arc::clone(&shared_check); |
| 563 | let job = Arc::clone(&job_check); |
| 564 | async move { job.checked(&shared).await.map_err(ToolError::not_available) } |
| 565 | }, |
| 566 | Arc::clone(&github.output), |
| 567 | ) |
| 568 | .await; |
| 569 | let projection = match result.as_ref() { |
| 570 | Ok(outcome) => Ok( |
| 571 | json!({"kind":"stock_adapter","operation":"github_result","input":&outcome.projection}), |
| 572 | ), |
| 573 | Err(_) => Err("GitHub Core operation failed".to_string()), |
| 574 | }; |
| 575 | github.output.lock().expect("GitHub result lock").result = Some(result); |
| 576 | // A Core fault never crosses into the Host diagnostic/projection. |
| 577 | let projection = projection?; |
| 578 | job.checked(&shared).await?; |
| 579 | return Ok(projection); |
| 580 | } |
| 581 | if let Some(ocr) = launch.ocr.take() { |
| 582 | match ocr { |
| 583 | OcrLaunch::Native { input, output } => { |
| 584 | let worker = Arc::clone(&job); |
| 585 | let shared_worker = Arc::clone(&shared); |
| 586 | #[cfg(test)] |
| 587 | let scope = crate::test_support::env_scope_ticket(); |
| 588 | let policy = activation::extension_host_policy_enabled(); |
| 589 | let step = shared |
| 590 | .engine_handle()? |
| 591 | .spawn_blocking(move || { |
| 592 | let _policy = activation::PolicyScope::propagate(policy); |
| 593 | #[cfg(test)] |
| 594 | let _scope = crate::test_support::join_env_scope(scope); |
| 595 | worker.check(&shared_worker)?; |
| 596 | // Retain the same job/quota until non-preemptible Vision FFI exits. |
| 597 | let step = input.native_step(); |
| 598 | worker.check(&shared_worker)?; |
| 599 | Ok::<_, String>(step) |
| 600 | }) |
| 601 | .await |
| 602 | .map_err(|_| "OCR Native worker failed".to_string())??; |
| 603 | job.checked(&shared).await?; |
| 604 | let mut projection = step.projection(); |
| 605 | if step.needs_tesseract() { |
| 606 | let mut continuation = job |
| 607 | .continuation |
| 608 | .lock() |
| 609 | .expect("execution continuation lock"); |
| 610 | job.check(&shared)?; |
| 611 | if continuation.is_some() { |
| 612 | return Err("OCR continuation was already issued".into()); |
| 613 | } |
| 614 | let remaining = job.expires.saturating_duration_since(Instant::now()); |
| 615 | if remaining.is_zero() { |
| 616 | return Err("OCR deadline exhausted".into()); |
| 617 | } |
| 618 | let mut target = job.target.clone(); |
| 619 | target["operation"] = json!("ocr_tesseract"); |
| 620 | target["captured_sha256"] = json!(step.digest()); |
| 621 | let ticket = shared.core_calls.tickets.mint(Grant { |
| 622 | kind: TicketKind::Execution, |
| 623 | tier: HostTier::Builtin, |
| 624 | host_generation: job.generation, |
| 625 | owner: job.owner.clone(), |
| 626 | method: "exec/redeem", |
| 627 | target: target.clone(), |
| 628 | ttl: remaining, |
| 629 | uses: 1, |
| 630 | }); |
| 631 | projection["next_ticket"] = json!(ticket.expose()); |
| 632 | *job.launch.lock().expect("execution launch lock") = Some(Launch { |
| 633 | github: None, |
| 634 | ocr: Some(OcrLaunch::Tesseract { step, output }), |
| 635 | pdf: None, |
| 636 | stock: None, |
| 637 | hook: None, |
| 638 | command: None, |
| 639 | input: Vec::new(), |
| 640 | }); |
| 641 | *continuation = Some(Continuation { ticket, target }); |
| 642 | } else { |
| 643 | *output.lock().expect("OCR result lock") = Some(step.finish()); |
| 644 | } |
| 645 | job.checked(&shared).await?; |
| 646 | return Ok(projection); |
| 647 | } |
| 648 | OcrLaunch::Tesseract { step, output } => { |
| 649 | let result = step |
| 650 | .tesseract( |
| 651 | Some(&job.cancel), |
| 652 | tokio::time::Instant::from_std(job.expires), |
| 653 | ) |
| 654 | .await; |
| 655 | let projection = result.projection(); |
| 656 | *output.lock().expect("OCR result lock") = Some(result); |
| 657 | job.checked(&shared).await?; |
| 658 | return Ok(projection); |
| 659 | } |
| 660 | } |
| 661 | } |
| 662 | if let Some(pdf) = launch.pdf.take() { |
| 663 | let result = crate::tools::pdf::run_pdf_driver(pdf.input, Some(&job.cancel)).await; |
| 664 | let projection = result.projection(); |
| 665 | *pdf.output.lock().expect("PDF result lock") = Some(result); |
| 666 | job.checked(&shared).await?; |
| 667 | return Ok(projection); |
| 668 | } |
| 669 | if let Some((operation, input)) = launch.stock.take() { |
| 670 | job.checked(&shared).await?; |
| 671 | return Ok(json!({"kind":"stock_adapter","operation":operation.as_str(),"input":input})); |
| 672 | } |
| 673 | if let Some(launch_hook) = launch.hook.take() { |
| 674 | let job_worker = Arc::clone(&job); |
| 675 | let shared_worker = Arc::clone(&shared); |
| 676 | let policy = activation::extension_host_policy_enabled(); |
| 677 | #[cfg(test)] |
| 678 | let env_scope = crate::test_support::env_scope_ticket(); |
| 679 | let result=shared.engine_handle()?.spawn_blocking(move || { |
| 680 | let _policy=activation::PolicyScope::propagate(policy); |
| 681 | #[cfg(test)] let _env_scope=crate::test_support::join_env_scope(env_scope); |
| 682 | job_worker.check(&shared_worker)?; |
| 683 | crate::hooks::authority::verify_hook(&launch_hook.hook)?; |
| 684 | if let Some(native)=launch_hook.hook.native_shell.as_ref() { |
| 685 | shared_worker.registry.lock().expect("registry lock").check_shell_hook(native)?; |
| 686 | } |
| 687 | let result=(launch_hook.run)(job_worker.cancel.clone()); |
| 688 | crate::hooks::authority::verify_hook(&launch_hook.hook)?; |
| 689 | if let Some(native)=launch_hook.hook.native_shell.as_ref() { |
| 690 | shared_worker.registry.lock().expect("registry lock").check_shell_hook(native)?; |
| 691 | } |
| 692 | job_worker.check(&shared_worker)?; |
| 693 | // ShellEnv stdout may contain credentials. It is retained only in Rust. |
| 694 | let projection=if launch_hook.hook.event==crate::hooks::HookEvent::ShellEnv { |
| 695 | json!({"kind":"hook","event":"shell_env","success":result.success,"exit_code":result.exit_code,"keys":crate::hooks::shell_env_keys(&result.stdout)}) |
| 696 | } else if launch_hook.hook.native_shell.is_some() { |
| 697 | json!({"kind":"hook","event":launch_hook.hook.event.as_str(),"success":result.success,"exit_code":result.exit_code,"stdout":result.stdout,"stderr":result.stderr}) |
| 698 | } else { |
| 699 | json!({"kind":"hook","event":launch_hook.hook.event.as_str(),"success":result.success,"exit_code":result.exit_code}) |
| 700 | }; |
| 701 | *job_worker.hook_result.lock().expect("hook result lock")=Some(result); |
| 702 | Ok::<Value,String>(projection) |
| 703 | }).await.map_err(|_|"hook worker was lost".to_string())??; |
| 704 | job.checked(&shared).await?; |
| 705 | return Ok(result); |
| 706 | } |
| 707 | let command = launch |
| 708 | .command |
| 709 | .as_mut() |
| 710 | .ok_or("execution command is missing")?; |
| 711 | // The selected Builtin never sees the environment, argv, input, cwd or decision. |
| 712 | crate::child_env::apply_to_tokio_command(command, std::iter::empty::<(&str, &str)>()); |
| 713 | job.check(&shared)?; |
| 714 | let stop = async { |
| 715 | let remaining = job.expires.saturating_duration_since(Instant::now()); |
| 716 | tokio::select! { _ = job.cancel.cancelled() => {}, _ = tokio::time::sleep(remaining) => {}, } |
| 717 | }; |
| 718 | let output = crate::process_tree::contained_output_with_input_bounded( |
| 719 | command, |
| 720 | launch.input, |
| 721 | MAX_STDOUT, |
| 722 | MAX_STDERR, |
| 723 | stop, |
| 724 | ) |
| 725 | .await |
| 726 | .map_err(|_| "execution spawn or output collection failed".to_string())?; |
| 727 | if output.stopped { |
| 728 | return Err("execution cancelled or timed out".into()); |
| 729 | } |
| 730 | job.checked(&shared).await?; |
| 731 | Ok( |
| 732 | json!({"success":output.output.status.success(),"stdout":String::from_utf8_lossy(&output.output.stdout),"stderr":String::from_utf8_lossy(&output.output.stderr)}), |
| 733 | ) |
| 734 | } |
| 735 | |
| 736 | impl ExtensionHostManager { |
| 737 | /// Same weighted quota for the nonpreemptible Core capture worker. |
| 738 | pub(crate) fn admit_review_capture(&self) -> Result<OwnedSemaphorePermit, ToolError> { |
| 739 | self.shared |
| 740 | .execution_broker |
| 741 | .admit(StockOperation::ReviewPassPrompt.slots()) |
| 742 | } |
| 743 | |
| 744 | /// Called only by the already planned/gated Script and Command ToolSpecs. |
| 745 | pub(crate) async fn execute_script( |
| 746 | &self, |
| 747 | command: tokio::process::Command, |
| 748 | input: Value, |
| 749 | context: &ToolContext, |
| 750 | ) -> Result<ToolResult, ToolError> { |
| 751 | let input = serde_json::to_vec(&input) |
| 752 | .map_err(|_| ToolError::invalid_input("script input is not JSON"))?; |
| 753 | if input.len() > MAX_INPUT { |
| 754 | return Err(ToolError::invalid_input("script input exceeds 1 MiB")); |
| 755 | } |
| 756 | self.execute_launch( |
| 757 | Launch { |
| 758 | github: None, |
| 759 | ocr: None, |
| 760 | pdf: None, |
| 761 | command: Some(command), |
| 762 | input, |
| 763 | hook: None, |
| 764 | stock: None, |
| 765 | }, |
| 766 | context, |
| 767 | DEADLINE, |
| 768 | None, |
| 769 | ) |
| 770 | .await |
| 771 | } |
| 772 | |
| 773 | /// Captured data only. Rust has already performed the planned network/file access |
| 774 | /// and parser checks. No URL, command, environment or writable handle is granted. |
| 775 | pub(crate) async fn execute_stock( |
| 776 | &self, |
| 777 | operation: StockOperation, |
| 778 | input: Value, |
| 779 | context: &ToolContext, |
| 780 | budget: Duration, |
| 781 | ) -> Result<ToolResult, ToolError> { |
| 782 | if operation.is_review() |
| 783 | && !context |
| 784 | .features |
| 785 | .enabled(crate::features::Feature::ReviewHost) |
| 786 | { |
| 787 | return Err(ToolError::not_available( |
| 788 | "Review Host backend is not selected", |
| 789 | )); |
| 790 | } |
| 791 | // Admit before constructing the encoded snapshot, and carry this exact |
| 792 | // permit into the actual job. Concurrent validation cannot outrun quota. |
| 793 | let admission = self.shared.execution_broker.admit(operation.slots())?; |
| 794 | check_stock_snapshot(operation, &input)?; |
| 795 | // Adopt the existing Engine scheduler; this never creates a runtime. |
| 796 | if let Ok(handle) = tokio::runtime::Handle::try_current() { |
| 797 | self.bind_engine_handle(handle); |
| 798 | } |
| 799 | self.execute_launch( |
| 800 | Launch { |
| 801 | github: None, |
| 802 | ocr: None, |
| 803 | pdf: None, |
| 804 | command: None, |
| 805 | input: Vec::new(), |
| 806 | hook: None, |
| 807 | stock: Some((operation, input)), |
| 808 | }, |
| 809 | context, |
| 810 | budget, |
| 811 | Some(admission), |
| 812 | ) |
| 813 | .await |
| 814 | } |
| 815 | |
| 816 | /// Only the canonical already-planned ToolSpec can capture a write operation. |
| 817 | pub(crate) async fn execute_github( |
| 818 | &self, |
| 819 | request: crate::tools::github::host::Request, |
| 820 | context: &ToolContext, |
| 821 | ) -> Result<ToolResult, ToolError> { |
| 822 | if !context |
| 823 | .features |
| 824 | .enabled(crate::features::Feature::GithubHost) |
| 825 | { |
| 826 | return Err(ToolError::not_available( |
| 827 | "GitHub Host backend is not selected", |
| 828 | )); |
| 829 | } |
| 830 | if let Ok(handle) = tokio::runtime::Handle::try_current() { |
| 831 | self.bind_engine_handle(handle); |
| 832 | } |
| 833 | let caller = Caller::capture(context, &self.shared).map_err(ToolError::not_available)?; |
| 834 | let cancel = context |
| 835 | .cancel_token |
| 836 | .as_ref() |
| 837 | .map(CancellationToken::child_token) |
| 838 | .unwrap_or_default(); |
| 839 | let mut budget = DEADLINE; |
| 840 | if let Some(deadline) = context.turn_deadline { |
| 841 | budget = budget.min(deadline.saturating_duration_since(tokio::time::Instant::now())); |
| 842 | } |
| 843 | self.execute_github_captured(request, Some(context.clone()), caller, cancel, budget) |
| 844 | .await |
| 845 | } |
| 846 | /// Read-only operator path: actual app/session/agent identity, no synthetic tool approval. |
| 847 | pub(crate) async fn execute_github_review( |
| 848 | &self, |
| 849 | caller: crate::hooks::HookCaller, |
| 850 | id: String, |
| 851 | ) -> Result<ToolResult, ToolError> { |
| 852 | let session = caller |
| 853 | .session_id |
| 854 | .clone() |
| 855 | .filter(|id| !id.is_empty()) |
| 856 | .ok_or_else(|| { |
| 857 | ToolError::not_available("Feedback review requires the current session") |
| 858 | })?; |
| 859 | if let Ok(handle) = tokio::runtime::Handle::try_current() { |
| 860 | self.bind_engine_handle(handle); |
| 861 | } |
| 862 | let caller = |
| 863 | Caller::from_hook(caller, &self.shared, None).map_err(ToolError::not_available)?; |
| 864 | self.execute_github_captured( |
| 865 | crate::tools::github::host::Request::Read { |
| 866 | session, |
| 867 | id, |
| 868 | operator: true, |
| 869 | }, |
| 870 | None, |
| 871 | caller, |
| 872 | CancellationToken::new(), |
| 873 | Duration::from_secs(30), |
| 874 | ) |
| 875 | .await |
| 876 | } |
| 877 | async fn execute_github_captured( |
| 878 | &self, |
| 879 | request: crate::tools::github::host::Request, |
| 880 | context: Option<ToolContext>, |
| 881 | caller: Caller, |
| 882 | cancel: CancellationToken, |
| 883 | budget: Duration, |
| 884 | ) -> Result<ToolResult, ToolError> { |
| 885 | let output = Arc::new(Mutex::new(crate::tools::github::host::RunState::default())); |
| 886 | let result = self |
| 887 | .execute_launch_for_caller( |
| 888 | Launch { |
| 889 | github: Some(GithubLaunch { |
| 890 | request, |
| 891 | context, |
| 892 | output: Arc::clone(&output), |
| 893 | }), |
| 894 | ocr: None, |
| 895 | pdf: None, |
| 896 | stock: None, |
| 897 | hook: None, |
| 898 | command: None, |
| 899 | input: Vec::new(), |
| 900 | }, |
| 901 | caller, |
| 902 | cancel, |
| 903 | budget, |
| 904 | None, |
| 905 | ) |
| 906 | .await; |
| 907 | crate::tools::github::host::finish_run(&output, result) |
| 908 | } |
| 909 | |
| 910 | /// Only the planned/gated PDF consumers can construct this captured parser job. |
| 911 | pub(crate) async fn execute_pdf( |
| 912 | &self, |
| 913 | input: crate::tools::pdf::CapturedPdf, |
| 914 | context: &ToolContext, |
| 915 | ) -> Result<(crate::tools::pdf::PdfProcessOutcome, ToolResult), ToolError> { |
| 916 | if !context.features.enabled(crate::features::Feature::PdfHost) { |
| 917 | return Err(ToolError::not_available("PDF Host backend is not selected")); |
| 918 | } |
| 919 | if let Ok(handle) = tokio::runtime::Handle::try_current() { |
| 920 | self.bind_engine_handle(handle); |
| 921 | } |
| 922 | let mut budget = input.timeout; |
| 923 | if let Some(deadline) = context.turn_deadline { |
| 924 | budget = budget.min(deadline.saturating_duration_since(tokio::time::Instant::now())); |
| 925 | } |
| 926 | let output = Arc::new(Mutex::new(None)); |
| 927 | let result = self |
| 928 | .execute_launch( |
| 929 | Launch { |
| 930 | github: None, |
| 931 | ocr: None, |
| 932 | pdf: Some(PdfLaunch { |
| 933 | input, |
| 934 | output: Arc::clone(&output), |
| 935 | }), |
| 936 | stock: None, |
| 937 | hook: None, |
| 938 | command: None, |
| 939 | input: Vec::new(), |
| 940 | }, |
| 941 | context, |
| 942 | budget, |
| 943 | None, |
| 944 | ) |
| 945 | .await?; |
| 946 | let output = output |
| 947 | .lock() |
| 948 | .expect("PDF result lock") |
| 949 | .take() |
| 950 | .ok_or_else(|| ToolError::execution_failed("PDF driver has no captured output"))?; |
| 951 | Ok((output, result)) |
| 952 | } |
| 953 | |
| 954 | pub(crate) async fn execute_ocr( |
| 955 | &self, |
| 956 | input: crate::tools::image_ocr::CapturedOcr, |
| 957 | context: &ToolContext, |
| 958 | ) -> Result<(crate::tools::image_ocr::OcrOutcome, ToolResult), ToolError> { |
| 959 | if !context.features.enabled(crate::features::Feature::OcrHost) { |
| 960 | return Err(ToolError::not_available("OCR Host backend is not selected")); |
| 961 | } |
| 962 | if let Ok(handle) = tokio::runtime::Handle::try_current() { |
| 963 | self.bind_engine_handle(handle); |
| 964 | } |
| 965 | let budget = context |
| 966 | .turn_deadline |
| 967 | .map(|deadline| deadline.saturating_duration_since(tokio::time::Instant::now())) |
| 968 | .unwrap_or(DEADLINE) |
| 969 | .min(DEADLINE); |
| 970 | let output = Arc::new(Mutex::new(None)); |
| 971 | let result = self |
| 972 | .execute_launch( |
| 973 | Launch { |
| 974 | github: None, |
| 975 | ocr: Some(OcrLaunch::Native { |
| 976 | input, |
| 977 | output: Arc::clone(&output), |
| 978 | }), |
| 979 | pdf: None, |
| 980 | stock: None, |
| 981 | hook: None, |
| 982 | command: None, |
| 983 | input: Vec::new(), |
| 984 | }, |
| 985 | context, |
| 986 | budget, |
| 987 | None, |
| 988 | ) |
| 989 | .await; |
| 990 | let decision = match result { |
| 991 | Err(_) |
| 992 | if context |
| 993 | .cancel_token |
| 994 | .as_ref() |
| 995 | .is_some_and(CancellationToken::is_cancelled) => |
| 996 | { |
| 997 | return Err(ToolError::cancelled("Image OCR was cancelled")); |
| 998 | } |
| 999 | other => other?, |
| 1000 | }; |
| 1001 | let outcome = output |
| 1002 | .lock() |
| 1003 | .expect("OCR result lock") |
| 1004 | .take() |
| 1005 | .ok_or_else(|| ToolError::execution_failed("OCR driver has no captured output"))?; |
| 1006 | Ok((outcome, decision)) |
| 1007 | } |
| 1008 | |
| 1009 | async fn execute_launch( |
| 1010 | &self, |
| 1011 | launch: Launch, |
| 1012 | context: &ToolContext, |
| 1013 | budget: Duration, |
| 1014 | admission: Option<OwnedSemaphorePermit>, |
| 1015 | ) -> Result<ToolResult, ToolError> { |
| 1016 | let caller = Caller::capture(context, &self.shared).map_err(ToolError::not_available)?; |
| 1017 | let cancel = context |
| 1018 | .cancel_token |
| 1019 | .as_ref() |
| 1020 | .map(CancellationToken::child_token) |
| 1021 | .unwrap_or_default(); |
| 1022 | self.execute_launch_for_caller(launch, caller, cancel, budget, admission) |
| 1023 | .await |
| 1024 | } |
| 1025 | |
| 1026 | async fn execute_launch_for_caller( |
| 1027 | &self, |
| 1028 | launch: Launch, |
| 1029 | mut caller: Caller, |
| 1030 | cancel: CancellationToken, |
| 1031 | budget: Duration, |
| 1032 | admission: Option<OwnedSemaphorePermit>, |
| 1033 | ) -> Result<ToolResult, ToolError> { |
| 1034 | self.shared |
| 1035 | .engine_handle() |
| 1036 | .map_err(ToolError::not_available)?; |
| 1037 | let budget = budget.min(DEADLINE); |
| 1038 | if budget.is_zero() { |
| 1039 | return Err(ToolError::Timeout { seconds: 0 }); |
| 1040 | } |
| 1041 | let expires = Instant::now() + budget; |
| 1042 | let stock = launch.stock.is_some() |
| 1043 | || launch.pdf.is_some() |
| 1044 | || launch.ocr.is_some() |
| 1045 | || launch.github.is_some(); |
| 1046 | caller.requires_native_policy = !stock; |
| 1047 | caller |
| 1048 | .check(&self.shared) |
| 1049 | .map_err(ToolError::not_available)?; |
| 1050 | let operation = launch.stock.as_ref().map(|(operation, _)| *operation); |
| 1051 | let review = operation.is_some_and(StockOperation::is_review); |
| 1052 | let permit = match admission { |
| 1053 | Some(admission) => admission, |
| 1054 | None => self |
| 1055 | .shared |
| 1056 | .execution_broker |
| 1057 | .admit(operation.map_or(1, StockOperation::slots))?, |
| 1058 | }; |
| 1059 | let demand = if stock { |
| 1060 | &self.shared.stock_users |
| 1061 | } else { |
| 1062 | &self.shared.harness_users |
| 1063 | }; |
| 1064 | demand.fetch_add(1, Ordering::SeqCst); |
| 1065 | let _demand = Demand(demand); |
| 1066 | let ready = async { |
| 1067 | if stock { |
| 1068 | self.ensure_stock_builtin().await |
| 1069 | } else { |
| 1070 | self.ensure_harness_builtin().await |
| 1071 | } |
| 1072 | }; |
| 1073 | |
| 1074 | let deadline = tokio::time::Instant::from_std(expires); |
| 1075 | let (host, owner, generation) = tokio::select! { |
| 1076 | biased; |
| 1077 | _ = cancel.cancelled() => return Err(ToolError::not_available("execution cancelled")), |
| 1078 | ready = tokio::time::timeout_at(deadline, ready) => ready |
| 1079 | .map_err(|_| ToolError::Timeout { seconds: budget.as_secs().saturating_add(u64::from(budget.subsec_nanos() != 0)) })? |
| 1080 | .map_err(ToolError::not_available)?, |
| 1081 | }; |
| 1082 | caller |
| 1083 | .check(&self.shared) |
| 1084 | .map_err(ToolError::not_available)?; |
| 1085 | let id = uuid::Uuid::new_v4().to_string(); |
| 1086 | let mut target = caller.target(&id); |
| 1087 | if let Some(github) = launch.github.as_ref() { |
| 1088 | target["operation"] = json!("github_capture"); |
| 1089 | target["snapshot_sha256"] = json!(crate::hashing::sha256_hex( |
| 1090 | &serde_json::to_vec(&github.request) |
| 1091 | .map_err(|_| ToolError::invalid_input("GitHub operation is not JSON"))? |
| 1092 | )); |
| 1093 | } |
| 1094 | if let Some(OcrLaunch::Native { input, .. }) = launch.ocr.as_ref() { |
| 1095 | target["operation"] = json!("ocr_native"); |
| 1096 | target["captured_sha256"] = json!(input.digest()); |
| 1097 | } |
| 1098 | if let Some(pdf) = launch.pdf.as_ref() { |
| 1099 | target["operation"] = json!("pdf_extract"); |
| 1100 | target["captured_sha256"] = json!(pdf.input.digest()); |
| 1101 | } |
| 1102 | if let Some((operation, input)) = launch.stock.as_ref() { |
| 1103 | target["operation"] = json!(operation.as_str()); |
| 1104 | target["snapshot_sha256"] = json!(crate::hashing::sha256_hex( |
| 1105 | &serde_json::to_vec(input) |
| 1106 | .map_err(|_| ToolError::invalid_input("adapter snapshot is not JSON"))? |
| 1107 | )); |
| 1108 | } |
| 1109 | let remaining = expires.saturating_duration_since(Instant::now()); |
| 1110 | if remaining.is_zero() { |
| 1111 | return Err(ToolError::not_available("execution deadline exhausted")); |
| 1112 | } |
| 1113 | let ticket = self.shared.core_calls.tickets.mint(Grant { |
| 1114 | kind: TicketKind::Execution, |
| 1115 | tier: HostTier::Builtin, |
| 1116 | host_generation: generation, |
| 1117 | owner: owner.clone(), |
| 1118 | method: "exec/redeem", |
| 1119 | target: target.clone(), |
| 1120 | ttl: remaining, |
| 1121 | uses: 1, |
| 1122 | }); |
| 1123 | let job = Arc::new(Job { |
| 1124 | id: id.clone(), |
| 1125 | owner: owner.clone(), |
| 1126 | generation, |
| 1127 | caller, |
| 1128 | target, |
| 1129 | expires, |
| 1130 | cancel, |
| 1131 | launch: Mutex::new(Some(launch)), |
| 1132 | ticket: ticket.clone(), |
| 1133 | continuation: Mutex::new(None), |
| 1134 | hook_result: Mutex::new(None), |
| 1135 | _permit: permit, |
| 1136 | }); |
| 1137 | self.shared |
| 1138 | .execution_broker |
| 1139 | .jobs |
| 1140 | .lock() |
| 1141 | .expect("execution lock") |
| 1142 | .insert(id.clone(), Arc::clone(&job)); |
| 1143 | let _invocation = Invocation { |
| 1144 | shared: Arc::clone(&self.shared), |
| 1145 | job: Arc::clone(&job), |
| 1146 | }; |
| 1147 | let request = host.call( |
| 1148 | CoreRequest::HarnessRun(HarnessRunParams { |
| 1149 | owner, |
| 1150 | execution_id: id, |
| 1151 | ticket: ticket.expose().to_string(), |
| 1152 | // Round the peer timer up: Core owns the exact deadline and must |
| 1153 | // settle expiry before a rounded-down Host timer reports a runner fault. |
| 1154 | deadline_ms: remaining.as_nanos().div_ceil(1_000_000).max(1) as u64, |
| 1155 | hook: None, |
| 1156 | }), |
| 1157 | Some("host:harness".into()), |
| 1158 | ); |
| 1159 | let value = tokio::select! { |
| 1160 | biased; |
| 1161 | _ = job.cancel.cancelled() => return Err(ToolError::not_available("execution cancelled")), |
| 1162 | value = tokio::time::timeout_at(deadline, request) => value |
| 1163 | .map_err(|_| ToolError::Timeout { seconds: budget.as_secs().saturating_add(u64::from(budget.subsec_nanos() != 0)) })? |
| 1164 | .map_err(|_| ToolError::execution_failed("Builtin runner failed; no Rust fallback was attempted"))?, |
| 1165 | }; |
| 1166 | job.checked(&self.shared) |
| 1167 | .await |
| 1168 | .map_err(ToolError::not_available)?; |
| 1169 | if review |
| 1170 | && !review_envelope_fits(&value).map_err(|_| { |
| 1171 | ToolError::execution_failed( |
| 1172 | "Serialized review result exceeds the existing frame limit", |
| 1173 | ) |
| 1174 | })? |
| 1175 | { |
| 1176 | return Err(ToolError::execution_failed( |
| 1177 | "Serialized review result/envelope exceeds 16 MiB", |
| 1178 | )); |
| 1179 | } |
| 1180 | if stock |
| 1181 | && !review |
| 1182 | && serde_json::to_vec(&value) |
| 1183 | .map_err(|_| ToolError::execution_failed("Builtin result is not JSON"))? |
| 1184 | .len() |
| 1185 | > MAX_INPUT |
| 1186 | { |
| 1187 | return Err(ToolError::execution_failed("Builtin result exceeds 1 MiB")); |
| 1188 | } |
| 1189 | if value.get("ok").and_then(Value::as_bool) == Some(true) { |
| 1190 | serde_json::from_value( |
| 1191 | value |
| 1192 | .get("result") |
| 1193 | .cloned() |
| 1194 | .ok_or_else(|| ToolError::execution_failed("Builtin result missing"))?, |
| 1195 | ) |
| 1196 | .map_err(|_| ToolError::execution_failed("Builtin result malformed")) |
| 1197 | } else { |
| 1198 | Err(ToolError::execution_failed( |
| 1199 | value |
| 1200 | .get("error") |
| 1201 | .and_then(Value::as_str) |
| 1202 | .unwrap_or("Builtin result failed"), |
| 1203 | )) |
| 1204 | } |
| 1205 | } |
| 1206 | } |
| 1207 | |
| 1208 | impl ExtensionHostManager { |
| 1209 | pub(crate) async fn execute_hook<F>( |
| 1210 | &self, |
| 1211 | caller: crate::hooks::HookCaller, |
| 1212 | hook: crate::hooks::Hook, |
| 1213 | timeout: Duration, |
| 1214 | query: String, |
| 1215 | run: F, |
| 1216 | ) -> Result<crate::hooks::HookResult, String> |
| 1217 | where |
| 1218 | F: FnOnce(CancellationToken) -> crate::hooks::HookResult + Send + 'static, |
| 1219 | { |
| 1220 | self.shared.engine_handle()?; |
| 1221 | let deadline_ms = u64::try_from(timeout.as_millis()) |
| 1222 | .map_err(|_| "hook timeout exceeds the execution bound")?; |
| 1223 | if deadline_ms == 0 || deadline_ms > i32::MAX as u64 { |
| 1224 | return Err("hook timeout exceeds the execution bound".into()); |
| 1225 | } |
| 1226 | let caller = Caller::from_hook( |
| 1227 | caller, |
| 1228 | &self.shared, |
| 1229 | Some( |
| 1230 | hook.native_shell |
| 1231 | .as_ref() |
| 1232 | .map_or("", |native| native.owner.plugin_id.as_str()), |
| 1233 | ), |
| 1234 | )?; |
| 1235 | caller.check(&self.shared)?; |
| 1236 | let permit = Arc::clone(&self.shared.execution_broker.slots) |
| 1237 | .try_acquire_owned() |
| 1238 | .map_err(|_| "execution concurrency limit reached")?; |
| 1239 | self.shared.harness_users.fetch_add(1, Ordering::SeqCst); |
| 1240 | let _demand = Demand(&self.shared.harness_users); |
| 1241 | let (host, owner, generation) = self.ensure_harness_builtin().await?; |
| 1242 | caller.check(&self.shared)?; |
| 1243 | let id = uuid::Uuid::new_v4().to_string(); |
| 1244 | let metadata = super::protocol::HookDispatchWire { |
| 1245 | event: hook.event.as_str().into(), |
| 1246 | dialect: hook |
| 1247 | .native_shell |
| 1248 | .as_ref() |
| 1249 | .map_or("codewhale".into(), |n| n.dialect.clone()), |
| 1250 | point: hook |
| 1251 | .native_shell |
| 1252 | .as_ref() |
| 1253 | .map_or(hook.event.as_str().into(), |n| n.point.clone()), |
| 1254 | matcher: hook.native_shell.as_ref().and_then(|n| n.matcher.clone()), |
| 1255 | query, |
| 1256 | }; |
| 1257 | let target = json!({"caller":caller.target(&id),"hook":metadata}); |
| 1258 | let ticket = self.shared.core_calls.tickets.mint(Grant { |
| 1259 | kind: TicketKind::Execution, |
| 1260 | tier: HostTier::Builtin, |
| 1261 | host_generation: generation, |
| 1262 | owner: owner.clone(), |
| 1263 | method: "exec/redeem", |
| 1264 | target: target.clone(), |
| 1265 | ttl: timeout, |
| 1266 | uses: 1, |
| 1267 | }); |
| 1268 | let job = Arc::new(Job { |
| 1269 | id: id.clone(), |
| 1270 | owner: owner.clone(), |
| 1271 | generation, |
| 1272 | caller, |
| 1273 | target, |
| 1274 | expires: Instant::now() |
| 1275 | .checked_add(timeout) |
| 1276 | .ok_or("hook timeout is not representable")?, |
| 1277 | cancel: CancellationToken::new(), |
| 1278 | launch: Mutex::new(Some(Launch { |
| 1279 | github: None, |
| 1280 | ocr: None, |
| 1281 | pdf: None, |
| 1282 | stock: None, |
| 1283 | command: None, |
| 1284 | input: Vec::new(), |
| 1285 | hook: Some(HookLaunch { |
| 1286 | run: Box::new(run), |
| 1287 | hook: hook.clone(), |
| 1288 | }), |
| 1289 | })), |
| 1290 | ticket: ticket.clone(), |
| 1291 | continuation: Mutex::new(None), |
| 1292 | hook_result: Mutex::new(None), |
| 1293 | _permit: permit, |
| 1294 | }); |
| 1295 | self.shared |
| 1296 | .execution_broker |
| 1297 | .jobs |
| 1298 | .lock() |
| 1299 | .expect("execution lock") |
| 1300 | .insert(id.clone(), Arc::clone(&job)); |
| 1301 | let _invocation = Invocation { |
| 1302 | shared: Arc::clone(&self.shared), |
| 1303 | job: Arc::clone(&job), |
| 1304 | }; |
| 1305 | let value = host |
| 1306 | .call( |
| 1307 | CoreRequest::HarnessRun(HarnessRunParams { |
| 1308 | owner, |
| 1309 | execution_id: id, |
| 1310 | ticket: ticket.expose().into(), |
| 1311 | deadline_ms, |
| 1312 | hook: Some(metadata), |
| 1313 | }), |
| 1314 | Some("host:harness".into()), |
| 1315 | ) |
| 1316 | .await |
| 1317 | .map_err(|_| "Builtin hook runner failed; no legacy fallback was attempted")?; |
| 1318 | job.checked(&self.shared).await?; |
| 1319 | if value.get("hook_skipped").and_then(Value::as_bool) == Some(true) { |
| 1320 | if job.hook_result.lock().expect("hook result lock").is_some() { |
| 1321 | return Err("hook ran despite its nonmatching receipt".into()); |
| 1322 | } |
| 1323 | return Ok(crate::hooks::HookResult { |
| 1324 | name: hook.name, |
| 1325 | success: true, |
| 1326 | background: false, |
| 1327 | strict: false, |
| 1328 | exit_code: Some(0), |
| 1329 | stdout: String::new(), |
| 1330 | stderr: String::new(), |
| 1331 | duration: Duration::ZERO, |
| 1332 | error: None, |
| 1333 | }); |
| 1334 | } |
| 1335 | if value.get("hook_completed").and_then(Value::as_bool) != Some(true) { |
| 1336 | return Err("Builtin hook receipt is missing".into()); |
| 1337 | } |
| 1338 | let mut result = job |
| 1339 | .hook_result |
| 1340 | .lock() |
| 1341 | .expect("hook result lock") |
| 1342 | .take() |
| 1343 | .ok_or("hook process produced no receipt")?; |
| 1344 | // Only Native dialect codecs may project a bounded proposal; final steering stays Rust-owned. |
| 1345 | if hook.native_shell.is_some() |
| 1346 | && hook.event != crate::hooks::HookEvent::ShellEnv |
| 1347 | && let Some(stdout) = value.get("proposal").and_then(Value::as_str) |
| 1348 | { |
| 1349 | if stdout.len() > 64 * 1024 { |
| 1350 | return Err("hook proposal exceeds 64 KiB".into()); |
| 1351 | } |
| 1352 | result.stdout = stdout.into(); |
| 1353 | } |
| 1354 | Ok(result) |
| 1355 | } |
| 1356 | } |
| 1357 | |
| 1358 | #[cfg(test)] |
| 1359 | mod tests { |
| 1360 | use super::super::{ExtensionHostOptions, ticket::TicketTable}; |
| 1361 | use super::*; |
| 1362 | fn caller() -> Caller { |
| 1363 | Caller { |
| 1364 | workspace: std::path::PathBuf::from("."), |
| 1365 | plugins: None, |
| 1366 | session_id: Some("session".into()), |
| 1367 | agent_id: Some("agent".into()), |
| 1368 | origin_turn_id: Some("turn".into()), |
| 1369 | origin_call_id: Some("call".into()), |
| 1370 | authorities: Vec::new(), |
| 1371 | native_owners: Vec::new(), |
| 1372 | native_host_generation: None, |
| 1373 | requires_native_policy: true, |
| 1374 | } |
| 1375 | } |
| 1376 | fn owner() -> OwnerRef { |
| 1377 | OwnerRef { |
| 1378 | plugin_id: "host:harness".into(), |
| 1379 | generation: 3, |
| 1380 | owner_token: "test-owner".into(), |
| 1381 | } |
| 1382 | } |
| 1383 | fn grant(owner: OwnerRef, target: Value) -> Grant { |
| 1384 | Grant { |
| 1385 | kind: TicketKind::Execution, |
| 1386 | tier: HostTier::Builtin, |
| 1387 | host_generation: 4, |
| 1388 | owner, |
| 1389 | method: "exec/redeem", |
| 1390 | target, |
| 1391 | ttl: DEADLINE, |
| 1392 | uses: 1, |
| 1393 | } |
| 1394 | } |
| 1395 | #[test] |
| 1396 | fn review_purpose_bounds_use_the_encoded_frame_and_keep_other_helpers_at_one_mib() { |
| 1397 | let large = json!({"kind":"cli_diff","diff":"漢".repeat(400_000)+"FINAL_END"}); |
| 1398 | assert!(serde_json::to_vec(&large).unwrap().len() > MAX_INPUT); |
| 1399 | assert!(check_stock_snapshot(StockOperation::ReviewSourcePrompt, &large).is_ok()); |
| 1400 | assert!(check_stock_snapshot(StockOperation::SpeechOptions, &large).is_err()); |
| 1401 | let escaped = json!({"kind":"cli_diff","diff":"\0".repeat(3*1024*1024)}); |
| 1402 | assert!(escaped["diff"].as_str().unwrap().len() < MAX_REVIEW); |
| 1403 | assert!(check_stock_snapshot(StockOperation::ReviewSourcePrompt, &escaped).is_err()); |
| 1404 | let result = json!({"ok":true,"result":{"content":"\0".repeat(3*1024*1024),"success":true,"metadata":null}}); |
| 1405 | assert!(!review_envelope_fits(&result).unwrap()); |
| 1406 | let frame = super::super::protocol::encode_frame( |
| 1407 | &json!({"jsonrpc":"2.0","id":u64::MAX,"result":large}), |
| 1408 | ) |
| 1409 | .unwrap(); |
| 1410 | assert_eq!( |
| 1411 | review_envelope_fits(&large).unwrap(), |
| 1412 | frame.len() <= MAX_REVIEW |
| 1413 | ); |
| 1414 | } |
| 1415 | |
| 1416 | #[test] |
| 1417 | fn weighted_review_admission_reuses_all_thirty_two_execution_slots() { |
| 1418 | let broker = Broker::default(); |
| 1419 | let first = broker.admit(StockOperation::ReviewReport.slots()).unwrap(); |
| 1420 | let second = broker |
| 1421 | .admit(StockOperation::ReviewPassPrompt.slots()) |
| 1422 | .unwrap(); |
| 1423 | assert_eq!(broker.slots.available_permits(), 0); |
| 1424 | assert!(broker.admit(1).is_err()); |
| 1425 | assert!(broker.admit(16).is_err()); |
| 1426 | drop(first); |
| 1427 | assert_eq!(broker.slots.available_permits(), 16); |
| 1428 | let ordinary = broker.admit(1).unwrap(); |
| 1429 | assert!(broker.admit(16).is_err()); |
| 1430 | drop(second); |
| 1431 | drop(ordinary); |
| 1432 | assert_eq!(broker.slots.available_permits(), MAX_JOBS); |
| 1433 | } |
| 1434 | |
| 1435 | #[tokio::test(flavor = "current_thread")] |
| 1436 | async fn abandoned_review_capture_worker_retains_weighted_quota_until_completion() { |
| 1437 | let manager = ExtensionHostManager::new(ExtensionHostOptions::default()); |
| 1438 | let permit = manager.admit_review_capture().unwrap(); |
| 1439 | let occupied = manager.admit_review_capture().unwrap(); |
| 1440 | let (started_tx, started_rx) = tokio::sync::oneshot::channel(); |
| 1441 | let (finish_tx, finish_rx) = std::sync::mpsc::channel(); |
| 1442 | let waiter = tokio::spawn(crate::tools::github::host::report_worker( |
| 1443 | permit, |
| 1444 | move || { |
| 1445 | started_tx.send(()).unwrap(); |
| 1446 | finish_rx.recv().unwrap(); |
| 1447 | Ok(()) |
| 1448 | }, |
| 1449 | )); |
| 1450 | started_rx.await.unwrap(); |
| 1451 | waiter.abort(); |
| 1452 | assert!(waiter.await.unwrap_err().is_cancelled()); |
| 1453 | assert!(manager.admit_review_capture().is_err()); |
| 1454 | assert_eq!(manager.shared.execution_broker.slots.available_permits(), 0); |
| 1455 | drop(occupied); |
| 1456 | assert_eq!( |
| 1457 | manager.shared.execution_broker.slots.available_permits(), |
| 1458 | 16 |
| 1459 | ); |
| 1460 | finish_tx.send(()).unwrap(); |
| 1461 | tokio::time::timeout(Duration::from_secs(2), async { |
| 1462 | while manager.shared.execution_broker.slots.available_permits() != MAX_JOBS { |
| 1463 | tokio::task::yield_now().await; |
| 1464 | } |
| 1465 | }) |
| 1466 | .await |
| 1467 | .unwrap(); |
| 1468 | assert!(manager.admit_review_capture().is_ok()); |
| 1469 | } |
| 1470 | |
| 1471 | #[tokio::test(flavor = "current_thread")] |
| 1472 | async fn actual_review_host_matches_frozen_core_cases_and_accepts_over_one_mib() { |
| 1473 | let _home = crate::test_support::SealedHome::new(); |
| 1474 | let _policy = crate::plugins::activation::TestPolicyGuard::extension_host(false); |
| 1475 | let Some(node) = super::super::tests::node_for_tests("review_captured_parity") else { |
| 1476 | return; |
| 1477 | }; |
| 1478 | let root = tempfile::tempdir().unwrap(); |
| 1479 | let manager = ExtensionHostManager::new(ExtensionHostOptions { |
| 1480 | runtime: crate::config::ExtensionHostRuntime::Node, |
| 1481 | node_override: Some(node), |
| 1482 | root: Some(root.path().join("host")), |
| 1483 | ..Default::default() |
| 1484 | }); |
| 1485 | let mut features = crate::features::Features::with_defaults(); |
| 1486 | features.enable(crate::features::Feature::ReviewHost); |
| 1487 | let context = ToolContext::new(root.path()).with_features(features); |
| 1488 | let fixture: Value = |
| 1489 | serde_json::from_str(include_str!("../../tests/fixtures/review-host-parity.json")) |
| 1490 | .unwrap(); |
| 1491 | for case in fixture["cases"].as_array().unwrap() { |
| 1492 | let operation = match case["operation"].as_str().unwrap() { |
| 1493 | "review_source_prompt" => StockOperation::ReviewSourcePrompt, |
| 1494 | "review_interactive_pr" => StockOperation::ReviewInteractivePr, |
| 1495 | "review_report" => StockOperation::ReviewReport, |
| 1496 | _ => panic!("unadmitted fixture operation"), |
| 1497 | }; |
| 1498 | let result = manager |
| 1499 | .execute_stock( |
| 1500 | operation, |
| 1501 | case["input"].clone(), |
| 1502 | &context, |
| 1503 | Duration::from_secs(10), |
| 1504 | ) |
| 1505 | .await |
| 1506 | .unwrap(); |
| 1507 | assert_eq!(result.content, case["content"].as_str().unwrap()); |
| 1508 | assert!(result.success); |
| 1509 | assert!(result.metadata.is_none()); |
| 1510 | } |
| 1511 | let input = json!({"kind":"cli_diff","diff":"漢".repeat(400_000)+"FINAL_END"}); |
| 1512 | let result = manager |
| 1513 | .execute_stock( |
| 1514 | StockOperation::ReviewSourcePrompt, |
| 1515 | input, |
| 1516 | &context, |
| 1517 | Duration::from_secs(10), |
| 1518 | ) |
| 1519 | .await |
| 1520 | .unwrap(); |
| 1521 | assert!(result.content.ends_with("FINAL_END\n\nEnd of diff.")); |
| 1522 | let cancelled = CancellationToken::new(); |
| 1523 | cancelled.cancel(); |
| 1524 | let error = manager |
| 1525 | .execute_stock( |
| 1526 | StockOperation::ReviewReport, |
| 1527 | json!({"review":null,"output":"late","posted":false}), |
| 1528 | &context.with_cancel_token(cancelled), |
| 1529 | Duration::from_secs(10), |
| 1530 | ) |
| 1531 | .await |
| 1532 | .unwrap_err(); |
| 1533 | assert!(matches!( |
| 1534 | error, |
| 1535 | ToolError::NotAvailable { .. } | ToolError::Cancelled { .. } |
| 1536 | )); |
| 1537 | manager.shutdown().await; |
| 1538 | } |
| 1539 | |
| 1540 | #[test] |
| 1541 | fn execution_ticket_cannot_redeem_mcp_or_replay_or_change_caller() { |
| 1542 | let table = TicketTable::default(); |
| 1543 | let owner = owner(); |
| 1544 | let caller = caller(); |
| 1545 | let target = caller.target("exact"); |
| 1546 | let ticket = table.mint(grant(owner.clone(), target.clone())); |
| 1547 | let mut present = Presented { |
| 1548 | ticket: ticket.expose(), |
| 1549 | kind: TicketKind::McpOperation, |
| 1550 | tier: HostTier::Builtin, |
| 1551 | host_generation: 4, |
| 1552 | owner: &owner, |
| 1553 | method: "exec/redeem", |
| 1554 | target: Some(&target), |
| 1555 | }; |
| 1556 | assert!(table.redeem(&present).is_err()); |
| 1557 | present.kind = TicketKind::Execution; |
| 1558 | let wrong = json!({"execution_id":"exact","session_id":"another"}); |
| 1559 | present.target = Some(&wrong); |
| 1560 | assert!(table.redeem(&present).is_err()); |
| 1561 | present.target = Some(&target); |
| 1562 | assert!(table.redeem(&present).is_ok()); |
| 1563 | assert!(table.redeem(&present).is_err()); |
| 1564 | } |
| 1565 | #[test] |
| 1566 | fn execution_caller_checks_exact_session_agent_and_selected_view() { |
| 1567 | let _policy = activation::PolicyScope::propagate(true); |
| 1568 | let manager = Arc::new(ExtensionHostManager::new(ExtensionHostOptions::default())); |
| 1569 | let attachment = manager.attach(Arc::new(PluginRegistry::empty(std::path::Path::new(".")))); |
| 1570 | attachment.set_identity(Some("session".into()), Some("agent".into())); |
| 1571 | let mut caller = caller(); |
| 1572 | caller.plugins = Some(attachment.plugin_view()); |
| 1573 | assert!(caller.check(&manager.shared).is_ok()); |
| 1574 | caller.session_id = Some("other".into()); |
| 1575 | assert!(caller.check(&manager.shared).is_err()); |
| 1576 | caller.session_id = Some("session".into()); |
| 1577 | caller.agent_id = None; |
| 1578 | assert!(caller.check(&manager.shared).is_err()); |
| 1579 | caller.agent_id = Some("agent".into()); |
| 1580 | drop(attachment); |
| 1581 | assert!(caller.check(&manager.shared).is_err()); |
| 1582 | } |
| 1583 | #[tokio::test(flavor = "current_thread")] |
| 1584 | async fn actual_github_broker_returns_private_core_permission_fault() { |
| 1585 | use crate::network_policy::{DecisionToml, NetworkPolicy, NetworkPolicyDecider}; |
| 1586 | let _home = crate::test_support::SealedHome::new(); |
| 1587 | let _policy = crate::plugins::activation::TestPolicyGuard::extension_host(false); |
| 1588 | let Some(node) = super::super::tests::node_for_tests("github_private_core_fault") else { |
| 1589 | return; |
| 1590 | }; |
| 1591 | let temp = tempfile::tempdir().unwrap(); |
| 1592 | let _repo = crate::test_support::EnvVarGuard::set("GH_REPO", "owner/name"); |
| 1593 | let manager = Arc::new(ExtensionHostManager::new(ExtensionHostOptions { |
| 1594 | runtime: crate::config::ExtensionHostRuntime::Node, |
| 1595 | node_override: Some(node), |
| 1596 | root: Some(temp.path().join("host")), |
| 1597 | ..Default::default() |
| 1598 | })); |
| 1599 | let mut context = |
| 1600 | ToolContext::new(temp.path()).with_network_policy(NetworkPolicyDecider::new( |
| 1601 | NetworkPolicy { |
| 1602 | default: DecisionToml::Deny, |
| 1603 | ..Default::default() |
| 1604 | }, |
| 1605 | None, |
| 1606 | )); |
| 1607 | context |
| 1608 | .features |
| 1609 | .enable(crate::features::Feature::GithubHost); |
| 1610 | let error = manager |
| 1611 | .execute_github( |
| 1612 | crate::tools::github::host::Request::Comment { |
| 1613 | target: "issue".into(), |
| 1614 | number: 1, |
| 1615 | body: "never sent".into(), |
| 1616 | dry: false, |
| 1617 | }, |
| 1618 | &context, |
| 1619 | ) |
| 1620 | .await |
| 1621 | .unwrap_err(); |
| 1622 | assert!(matches!(error, ToolError::PermissionDenied { .. })); |
| 1623 | assert!( |
| 1624 | error |
| 1625 | .to_string() |
| 1626 | .contains("blocked or awaiting network approval") |
| 1627 | ); |
| 1628 | assert!(!error.to_string().contains("Builtin runner failed")); |
| 1629 | assert!(!error.to_string().contains("never sent")); |
| 1630 | manager.shutdown().await; |
| 1631 | } |
| 1632 | |
| 1633 | #[cfg(unix)] |
| 1634 | #[tokio::test(flavor = "current_thread")] |
| 1635 | async fn actual_ocr_native_cancel_retains_job_quota_until_framework_worker_exits() { |
| 1636 | let _home = crate::test_support::SealedHome::new(); |
| 1637 | let _policy = crate::plugins::activation::TestPolicyGuard::extension_host(false); |
| 1638 | let Some(node) = super::super::tests::node_for_tests("ocr_native_quota") else { |
| 1639 | return; |
| 1640 | }; |
| 1641 | let root = tempfile::tempdir().unwrap(); |
| 1642 | let manager = Arc::new(ExtensionHostManager::new(ExtensionHostOptions { |
| 1643 | runtime: crate::config::ExtensionHostRuntime::Node, |
| 1644 | node_override: Some(node), |
| 1645 | root: Some(root.path().join("host")), |
| 1646 | ..ExtensionHostOptions::default() |
| 1647 | })); |
| 1648 | let started = Arc::new(tokio::sync::Notify::new()); |
| 1649 | let release = Arc::new((Mutex::new(false), std::sync::Condvar::new())); |
| 1650 | let announce = Arc::clone(&started); |
| 1651 | let latch = Arc::clone(&release); |
| 1652 | let input = crate::tools::image_ocr::CapturedOcr::for_test( |
| 1653 | root.path(), |
| 1654 | Arc::new(move |_| { |
| 1655 | announce.notify_one(); |
| 1656 | let (lock, wake) = &*latch; |
| 1657 | let mut ready = lock.lock().unwrap(); |
| 1658 | while !*ready { |
| 1659 | ready = wake.wait(ready).unwrap(); |
| 1660 | } |
| 1661 | Ok(Some("discarded late Native text".into())) |
| 1662 | }), |
| 1663 | None, |
| 1664 | ); |
| 1665 | let mut flags = crate::features::Features::with_defaults(); |
| 1666 | flags.enable(crate::features::Feature::OcrHost); |
| 1667 | let cancel = CancellationToken::new(); |
| 1668 | let context = ToolContext::new(root.path()) |
| 1669 | .with_features(flags) |
| 1670 | .with_cancel_token(cancel.clone()); |
| 1671 | let future = manager.execute_ocr(input, &context); |
| 1672 | tokio::pin!(future); |
| 1673 | tokio::select! {result=&mut future=>panic!("OCR settled before its blocking worker: {:?}",result.err()),_=started.notified()=>{},} |
| 1674 | cancel.cancel(); |
| 1675 | assert!(matches!(future.await, Err(ToolError::Cancelled { .. }))); |
| 1676 | assert_eq!( |
| 1677 | manager.shared.execution_broker.slots.available_permits(), |
| 1678 | MAX_JOBS - 1 |
| 1679 | ); |
| 1680 | let (lock, wake) = &*release; |
| 1681 | *lock.lock().unwrap() = true; |
| 1682 | wake.notify_all(); |
| 1683 | tokio::time::timeout(Duration::from_secs(5), async { |
| 1684 | while manager.shared.execution_broker.slots.available_permits() != MAX_JOBS { |
| 1685 | tokio::time::sleep(Duration::from_millis(5)).await; |
| 1686 | } |
| 1687 | }) |
| 1688 | .await |
| 1689 | .unwrap(); |
| 1690 | manager.shutdown().await; |
| 1691 | } |
| 1692 | |
| 1693 | #[cfg(unix)] |
| 1694 | #[tokio::test(flavor = "current_thread")] |
| 1695 | async fn actual_ocr_broker_decodes_exact_ordered_grants_and_refuses_replay_or_withdrawal() { |
| 1696 | use std::os::unix::fs::PermissionsExt; |
| 1697 | let _home = crate::test_support::SealedHome::new(); |
| 1698 | let _policy = crate::plugins::activation::TestPolicyGuard::extension_host(false); |
| 1699 | let Some(node) = super::super::tests::node_for_tests("ocr_exact_grants") else { |
| 1700 | return; |
| 1701 | }; |
| 1702 | let root = tempfile::tempdir().unwrap(); |
| 1703 | let marker = root.path().join("written"); |
| 1704 | let binary = root.path().join("tesseract"); |
| 1705 | std::fs::write( |
| 1706 | &binary, |
| 1707 | format!( |
| 1708 | "#!/bin/sh\nprintf x >> '{}'; printf 'private extracted text\\n'\n", |
| 1709 | marker.display() |
| 1710 | ), |
| 1711 | ) |
| 1712 | .unwrap(); |
| 1713 | std::fs::set_permissions(&binary, std::fs::Permissions::from_mode(0o700)).unwrap(); |
| 1714 | let manager = Arc::new(ExtensionHostManager::new(ExtensionHostOptions { |
| 1715 | runtime: crate::config::ExtensionHostRuntime::Node, |
| 1716 | node_override: Some(node), |
| 1717 | root: Some(root.path().join("host")), |
| 1718 | ..ExtensionHostOptions::default() |
| 1719 | })); |
| 1720 | manager.bind_engine_handle(tokio::runtime::Handle::current()); |
| 1721 | manager.shared.stock_users.fetch_add(1, Ordering::SeqCst); |
| 1722 | let _demand = Demand(&manager.shared.stock_users); |
| 1723 | let (_host, owner, generation) = manager.ensure_stock_builtin().await.unwrap(); |
| 1724 | let context = ToolContext::new(root.path()); |
| 1725 | let make_job = |id: &str| { |
| 1726 | let input = crate::tools::image_ocr::CapturedOcr::for_test( |
| 1727 | root.path(), |
| 1728 | Arc::new(|_| Err(ToolError::execution_failed("Native refusal"))), |
| 1729 | Some(binary.clone().into_os_string()), |
| 1730 | ); |
| 1731 | let output = Arc::new(Mutex::new(None)); |
| 1732 | let mut caller = Caller::capture(&context, &manager.shared).unwrap(); |
| 1733 | caller.requires_native_policy = false; |
| 1734 | let mut target = caller.target(id); |
| 1735 | target["operation"] = json!("ocr_native"); |
| 1736 | target["captured_sha256"] = json!(input.digest()); |
| 1737 | let ticket = manager.shared.core_calls.tickets.mint(Grant { |
| 1738 | kind: TicketKind::Execution, |
| 1739 | tier: HostTier::Builtin, |
| 1740 | host_generation: generation, |
| 1741 | owner: owner.clone(), |
| 1742 | method: "exec/redeem", |
| 1743 | target: target.clone(), |
| 1744 | ttl: DEADLINE, |
| 1745 | uses: 1, |
| 1746 | }); |
| 1747 | let job = Arc::new(Job { |
| 1748 | id: id.into(), |
| 1749 | owner: owner.clone(), |
| 1750 | generation, |
| 1751 | caller, |
| 1752 | target, |
| 1753 | expires: Instant::now() + DEADLINE, |
| 1754 | cancel: CancellationToken::new(), |
| 1755 | launch: Mutex::new(Some(Launch { |
| 1756 | github: None, |
| 1757 | ocr: Some(OcrLaunch::Native { |
| 1758 | input, |
| 1759 | output: Arc::clone(&output), |
| 1760 | }), |
| 1761 | pdf: None, |
| 1762 | stock: None, |
| 1763 | hook: None, |
| 1764 | command: None, |
| 1765 | input: Vec::new(), |
| 1766 | })), |
| 1767 | ticket: ticket.clone(), |
| 1768 | continuation: Mutex::new(None), |
| 1769 | hook_result: Mutex::new(None), |
| 1770 | _permit: Arc::clone(&manager.shared.execution_broker.slots) |
| 1771 | .try_acquire_owned() |
| 1772 | .unwrap(), |
| 1773 | }); |
| 1774 | manager |
| 1775 | .shared |
| 1776 | .execution_broker |
| 1777 | .jobs |
| 1778 | .lock() |
| 1779 | .unwrap() |
| 1780 | .insert(id.into(), Arc::clone(&job)); |
| 1781 | let invocation = Invocation { |
| 1782 | shared: Arc::clone(&manager.shared), |
| 1783 | job: Arc::clone(&job), |
| 1784 | }; |
| 1785 | (job, invocation, output) |
| 1786 | }; |
| 1787 | let (job, invocation, output) = make_job("ordered-ocr"); |
| 1788 | let decoded = |id: &str, ticket: &str| { |
| 1789 | serde_json::from_value::<ExecutionRedeemParams>( |
| 1790 | json!({"owner":owner,"execution_id":id,"ticket":ticket}), |
| 1791 | ) |
| 1792 | .unwrap() |
| 1793 | }; |
| 1794 | let cx = || HostRequestContext::for_test(99).0; |
| 1795 | let first = decoded(&job.id, job.ticket.expose()); |
| 1796 | let mut wrong_owner = first.clone(); |
| 1797 | wrong_owner.owner.plugin_id = "native:forged".into(); |
| 1798 | assert!( |
| 1799 | manager |
| 1800 | .shared |
| 1801 | .execution_broker |
| 1802 | .serve(&manager.shared, generation, wrong_owner, cx()) |
| 1803 | .await |
| 1804 | .is_err() |
| 1805 | ); |
| 1806 | assert!( |
| 1807 | manager |
| 1808 | .shared |
| 1809 | .execution_broker |
| 1810 | .serve(&manager.shared, generation + 1, first.clone(), cx()) |
| 1811 | .await |
| 1812 | .is_err() |
| 1813 | ); |
| 1814 | assert!( |
| 1815 | manager |
| 1816 | .shared |
| 1817 | .execution_broker |
| 1818 | .serve( |
| 1819 | &manager.shared, |
| 1820 | generation, |
| 1821 | decoded("other-ocr", job.ticket.expose()), |
| 1822 | cx() |
| 1823 | ) |
| 1824 | .await |
| 1825 | .is_err() |
| 1826 | ); |
| 1827 | assert!(!marker.exists()); |
| 1828 | let native = manager |
| 1829 | .shared |
| 1830 | .execution_broker |
| 1831 | .serve(&manager.shared, generation, first.clone(), cx()) |
| 1832 | .await |
| 1833 | .unwrap(); |
| 1834 | let next = native["next_ticket"].as_str().unwrap().to_string(); |
| 1835 | assert_eq!(native["status"], "error"); |
| 1836 | assert!(!marker.exists()); |
| 1837 | // The first grant is already spent and never authorizes the later write. |
| 1838 | assert!( |
| 1839 | manager |
| 1840 | .shared |
| 1841 | .execution_broker |
| 1842 | .serve(&manager.shared, generation, first, cx()) |
| 1843 | .await |
| 1844 | .is_err() |
| 1845 | ); |
| 1846 | let target = job |
| 1847 | .continuation |
| 1848 | .lock() |
| 1849 | .unwrap() |
| 1850 | .as_ref() |
| 1851 | .unwrap() |
| 1852 | .target |
| 1853 | .clone(); |
| 1854 | assert_eq!(target["operation"], "ocr_tesseract"); |
| 1855 | assert_ne!(target["captured_sha256"], job.target["captured_sha256"]); |
| 1856 | let wrong_kind = manager.shared.core_calls.tickets.mint(Grant { |
| 1857 | kind: TicketKind::McpOperation, |
| 1858 | tier: HostTier::Builtin, |
| 1859 | host_generation: generation, |
| 1860 | owner: owner.clone(), |
| 1861 | method: "exec/redeem", |
| 1862 | target, |
| 1863 | ttl: DEADLINE, |
| 1864 | uses: 1, |
| 1865 | }); |
| 1866 | assert!( |
| 1867 | manager |
| 1868 | .shared |
| 1869 | .execution_broker |
| 1870 | .serve( |
| 1871 | &manager.shared, |
| 1872 | generation, |
| 1873 | decoded(&job.id, wrong_kind.expose()), |
| 1874 | cx() |
| 1875 | ) |
| 1876 | .await |
| 1877 | .is_err() |
| 1878 | ); |
| 1879 | assert!(!marker.exists()); |
| 1880 | let mut altered_target = job |
| 1881 | .continuation |
| 1882 | .lock() |
| 1883 | .unwrap() |
| 1884 | .as_ref() |
| 1885 | .unwrap() |
| 1886 | .target |
| 1887 | .clone(); |
| 1888 | altered_target["captured_sha256"] = json!("changed-command-or-image"); |
| 1889 | let mismatch = manager.shared.core_calls.tickets.mint(Grant { |
| 1890 | kind: TicketKind::Execution, |
| 1891 | tier: HostTier::Builtin, |
| 1892 | host_generation: generation, |
| 1893 | owner: owner.clone(), |
| 1894 | method: "exec/redeem", |
| 1895 | target: altered_target, |
| 1896 | ttl: DEADLINE, |
| 1897 | uses: 1, |
| 1898 | }); |
| 1899 | assert!( |
| 1900 | manager |
| 1901 | .shared |
| 1902 | .execution_broker |
| 1903 | .serve( |
| 1904 | &manager.shared, |
| 1905 | generation, |
| 1906 | decoded(&job.id, mismatch.expose()), |
| 1907 | cx() |
| 1908 | ) |
| 1909 | .await |
| 1910 | .is_err() |
| 1911 | ); |
| 1912 | assert!( |
| 1913 | !marker.exists(), |
| 1914 | "parsed request must match its exact captured-operation digest before write" |
| 1915 | ); |
| 1916 | assert!(serde_json::from_value::<ExecutionRedeemParams>(json!({"owner":owner,"execution_id":job.id,"ticket":next,"path":"unreviewed-image.png"})).is_err()); |
| 1917 | let next = decoded(&job.id, &next); |
| 1918 | let projected = manager |
| 1919 | .shared |
| 1920 | .execution_broker |
| 1921 | .serve(&manager.shared, generation, next.clone(), cx()) |
| 1922 | .await |
| 1923 | .unwrap(); |
| 1924 | assert_eq!(projected["state"], "tesseract"); |
| 1925 | assert!(projected.get("stdout").is_none()); |
| 1926 | assert!(output.lock().unwrap().is_some()); |
| 1927 | assert_eq!(std::fs::read(&marker).unwrap(), b"x"); |
| 1928 | assert!( |
| 1929 | manager |
| 1930 | .shared |
| 1931 | .execution_broker |
| 1932 | .serve(&manager.shared, generation, next, cx()) |
| 1933 | .await |
| 1934 | .is_err() |
| 1935 | ); |
| 1936 | assert_eq!(std::fs::read(&marker).unwrap(), b"x"); |
| 1937 | drop(invocation); |
| 1938 | let (withdrawn, invocation, _) = make_job("withdrawn-ocr"); |
| 1939 | let native = manager |
| 1940 | .shared |
| 1941 | .execution_broker |
| 1942 | .serve( |
| 1943 | &manager.shared, |
| 1944 | generation, |
| 1945 | decoded(&withdrawn.id, withdrawn.ticket.expose()), |
| 1946 | cx(), |
| 1947 | ) |
| 1948 | .await |
| 1949 | .unwrap(); |
| 1950 | let next = decoded(&withdrawn.id, native["next_ticket"].as_str().unwrap()); |
| 1951 | manager.shared.execution_broker.revoke_owner("host:harness"); |
| 1952 | assert!( |
| 1953 | manager |
| 1954 | .shared |
| 1955 | .execution_broker |
| 1956 | .serve(&manager.shared, generation, next, cx()) |
| 1957 | .await |
| 1958 | .is_err() |
| 1959 | ); |
| 1960 | assert_eq!(std::fs::read(&marker).unwrap(), b"x"); |
| 1961 | drop(invocation); |
| 1962 | manager.shutdown().await; |
| 1963 | } |
| 1964 | |
| 1965 | #[tokio::test(flavor = "current_thread")] |
| 1966 | async fn abandoned_execution_worker_keeps_quota_until_cleanup_completes() { |
| 1967 | let manager = Arc::new(ExtensionHostManager::new(ExtensionHostOptions::default())); |
| 1968 | manager.bind_engine_handle(tokio::runtime::Handle::current()); |
| 1969 | let broker = &manager.shared.execution_broker; |
| 1970 | let owner = owner(); |
| 1971 | let caller = caller(); |
| 1972 | let target = caller.target("job"); |
| 1973 | let ticket = manager |
| 1974 | .shared |
| 1975 | .core_calls |
| 1976 | .tickets |
| 1977 | .mint(grant(owner.clone(), target.clone())); |
| 1978 | let job = Arc::new(Job { |
| 1979 | id: "job".into(), |
| 1980 | owner, |
| 1981 | generation: 4, |
| 1982 | caller, |
| 1983 | target, |
| 1984 | expires: Instant::now() + DEADLINE, |
| 1985 | cancel: CancellationToken::new(), |
| 1986 | launch: Mutex::new(None), |
| 1987 | ticket, |
| 1988 | continuation: Mutex::new(None), |
| 1989 | hook_result: Mutex::new(None), |
| 1990 | _permit: Arc::clone(&broker.slots).acquire_owned().await.unwrap(), |
| 1991 | }); |
| 1992 | broker |
| 1993 | .jobs |
| 1994 | .lock() |
| 1995 | .unwrap() |
| 1996 | .insert("job".into(), Arc::clone(&job)); |
| 1997 | let invocation = Invocation { |
| 1998 | shared: Arc::clone(&manager.shared), |
| 1999 | job: Arc::clone(&job), |
| 2000 | }; |
| 2001 | let (started_tx, started_rx) = tokio::sync::oneshot::channel(); |
| 2002 | let (finish_tx, finish_rx) = tokio::sync::oneshot::channel(); |
| 2003 | let (done_tx, done_rx) = tokio::sync::oneshot::channel(); |
| 2004 | let worker = dispatch_job( |
| 2005 | manager.shared.engine_handle().unwrap(), |
| 2006 | Arc::clone(&job), |
| 2007 | move |job| async move { |
| 2008 | started_tx.send(()).unwrap(); |
| 2009 | job.cancel.cancelled().await; |
| 2010 | finish_rx.await.unwrap(); |
| 2011 | done_tx.send(()).unwrap(); |
| 2012 | Ok(json!({})) |
| 2013 | }, |
| 2014 | ); |
| 2015 | started_rx.await.unwrap(); |
| 2016 | drop(worker); |
| 2017 | drop(invocation); |
| 2018 | drop(job); |
| 2019 | assert_eq!(broker.slots.available_permits(), MAX_JOBS - 1); |
| 2020 | finish_tx.send(()).unwrap(); |
| 2021 | done_rx.await.unwrap(); |
| 2022 | tokio::task::yield_now().await; |
| 2023 | assert_eq!(broker.slots.available_permits(), MAX_JOBS); |
| 2024 | } |
| 2025 | } |
| 2026 |