返回 CodeWhale
execution.rs
根目录 / crates / tui / src / extension_host / execution.rs
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, &current) {
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(&params.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: &params.ticket,
457 kind: TicketKind::Execution,
458 tier: HostTier::Builtin,
459 host_generation: generation,
460 owner: &params.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
2026 lines RUST