Skip to main content

lenso_kernel/
runtime.rs

1use super::{
2    AppAdmission, AppReadyGate, BTreeMap, CancellationToken, Cell, DiagnosticEvent,
3    DiagnosticShutdownOutcome, DiagnosticSource, DriverControl, DriverTask, Duration,
4    EventCapability, ExecutionAdapterCatalog, InvocationContext, LocalBoxFuture,
5    ManagedResourceScope, ManagedTask, ManagedTaskScope, NativeEventBindingTable,
6    NativeEventEndpointStateTable, NativeEventHandle, NativeRequestEndpoint, NativeRequestHandle,
7    NativeStreamBindingTable, NativeStreamEndpoint, NativeStreamEndpointStateTable,
8    NativeStreamHandle, PluginCriticality, PluginDependencies, PluginLifecycle, Rc, RefCell,
9    RequestAdmission, RequestCapability, RequestId, ResolvedAppPlan, RestartPolicy,
10    RuntimeDiagnostics, RuntimeFailure, ShutdownOutcome, StreamCapability,
11    begin_plugin_supervision, event, handle_supervision_schedule_failure, oneshot,
12    schedule_plugin_supervision, shutdown_native_plugins,
13};
14
15#[derive(Clone, Debug)]
16pub(super) struct NativeEndpointSnapshot {
17    pub(super) endpoint: Rc<dyn NativeRequestEndpoint>,
18    pub(super) generation: u64,
19    pub(super) cancellation: CancellationToken,
20}
21
22#[derive(Debug)]
23pub(super) struct NativeEndpointState {
24    pub(super) capability_id: &'static str,
25    pub(super) descriptor_version: &'static str,
26    pub(super) operations: &'static [&'static str],
27    pub(super) endpoint: RefCell<Option<Rc<dyn NativeRequestEndpoint>>>,
28    pub(super) generation: Cell<u64>,
29    pub(super) cancellation: RefCell<CancellationToken>,
30}
31
32#[derive(Clone, Debug)]
33pub(crate) struct NativeStreamEndpointSnapshot {
34    pub(crate) endpoint: Rc<dyn NativeStreamEndpoint>,
35    pub(crate) generation: u64,
36    pub(crate) cancellation: CancellationToken,
37}
38
39#[derive(Debug)]
40pub(crate) struct NativeStreamEndpointState {
41    pub(super) capability_id: &'static str,
42    pub(super) descriptor_version: &'static str,
43    pub(super) operations: &'static [&'static str],
44    pub(super) endpoint: RefCell<Option<Rc<dyn NativeStreamEndpoint>>>,
45    pub(super) generation: Cell<u64>,
46    pub(super) cancellation: RefCell<CancellationToken>,
47}
48
49impl NativeStreamEndpointState {
50    pub(crate) fn new(endpoint: Rc<dyn NativeStreamEndpoint>, generation: u64) -> Self {
51        Self {
52            capability_id: endpoint.capability_id(),
53            descriptor_version: endpoint.descriptor_version(),
54            operations: endpoint.operations(),
55            endpoint: RefCell::new(Some(endpoint)),
56            generation: Cell::new(generation),
57            cancellation: RefCell::new(CancellationToken::new()),
58        }
59    }
60
61    pub(crate) fn snapshot(&self) -> Option<NativeStreamEndpointSnapshot> {
62        self.endpoint
63            .borrow()
64            .clone()
65            .map(|endpoint| NativeStreamEndpointSnapshot {
66                endpoint,
67                generation: self.generation.get(),
68                cancellation: self.cancellation.borrow().clone(),
69            })
70    }
71
72    pub(crate) fn mark_unavailable(&self) {
73        self.cancellation.borrow().cancel();
74        self.endpoint.borrow_mut().take();
75    }
76
77    pub(crate) fn cancel(&self) {
78        self.cancellation.borrow().cancel();
79    }
80
81    pub(crate) fn install(&self, endpoint: Rc<dyn NativeStreamEndpoint>, generation: u64) {
82        self.generation.set(generation);
83        self.cancellation.replace(CancellationToken::new());
84        self.endpoint.replace(Some(endpoint));
85    }
86
87    pub(crate) fn is_current(&self, generation: u64) -> bool {
88        self.generation.get() == generation && self.endpoint.borrow().is_some()
89    }
90}
91
92impl NativeEndpointState {
93    pub(super) fn new(endpoint: Rc<dyn NativeRequestEndpoint>, generation: u64) -> Self {
94        Self {
95            capability_id: endpoint.capability_id(),
96            descriptor_version: endpoint.descriptor_version(),
97            operations: endpoint.operations(),
98            endpoint: RefCell::new(Some(endpoint)),
99            generation: Cell::new(generation),
100            cancellation: RefCell::new(CancellationToken::new()),
101        }
102    }
103
104    pub(super) fn snapshot(&self) -> Option<NativeEndpointSnapshot> {
105        self.endpoint
106            .borrow()
107            .clone()
108            .map(|endpoint| NativeEndpointSnapshot {
109                endpoint,
110                generation: self.generation.get(),
111                cancellation: self.cancellation.borrow().clone(),
112            })
113    }
114
115    pub(super) fn mark_unavailable(&self) {
116        self.cancellation.borrow().cancel();
117        self.endpoint.borrow_mut().take();
118    }
119
120    pub(super) fn cancel(&self) {
121        self.cancellation.borrow().cancel();
122    }
123
124    pub(super) fn install(&self, endpoint: Rc<dyn NativeRequestEndpoint>, generation: u64) {
125        self.generation.set(generation);
126        self.cancellation.replace(CancellationToken::new());
127        self.endpoint.replace(Some(endpoint));
128    }
129
130    pub(super) fn is_current(&self, generation: u64) -> bool {
131        self.generation.get() == generation && self.endpoint.borrow().is_some()
132    }
133}
134
135#[derive(Clone, Debug)]
136pub(super) struct NativeEndpointBinding {
137    pub(super) requirement_id: String,
138    pub(super) plugin_instance: String,
139    pub(super) state: Rc<NativeEndpointState>,
140    pub(super) admissions: BTreeMap<String, RequestAdmission>,
141}
142
143impl NativeEndpointBinding {
144    pub(super) fn admission(&self, operation: &str) -> Option<&RequestAdmission> {
145        self.admissions.get(operation)
146    }
147}
148
149#[derive(Clone, Debug)]
150pub(crate) struct NativeStreamEndpointBinding {
151    pub(super) requirement_id: String,
152    pub(crate) plugin_instance: String,
153    pub(crate) state: Rc<NativeStreamEndpointState>,
154    pub(super) admissions: BTreeMap<String, RequestAdmission>,
155}
156
157impl NativeStreamEndpointBinding {
158    pub(crate) fn admission(&self, operation: &str) -> Option<&RequestAdmission> {
159        self.admissions.get(operation)
160    }
161}
162
163#[derive(Debug)]
164pub(super) struct NativePluginGeneration {
165    pub(super) lifecycle: Rc<dyn PluginLifecycle>,
166    pub(super) tasks: ManagedTaskScope,
167    pub(super) resources: ManagedResourceScope,
168    pub(super) stop_attempted: bool,
169    pub(super) cleanup_timed_out: bool,
170}
171
172pub(super) enum GenerationPreparationFailure {
173    Lifecycle,
174    Cleanup { primary: RuntimeFailure },
175}
176
177#[derive(Debug)]
178pub(super) struct NativePluginRuntime {
179    pub(super) generation: RefCell<Option<NativePluginGeneration>>,
180}
181
182impl NativePluginRuntime {
183    pub(super) fn take_generation(&self) -> Option<NativePluginGeneration> {
184        self.generation.borrow_mut().take()
185    }
186
187    pub(super) fn install_generation(&self, generation: NativePluginGeneration) {
188        debug_assert!(self.generation.borrow().is_none());
189        self.generation.replace(Some(generation));
190    }
191
192    pub(super) fn generation_parts(
193        &self,
194    ) -> Option<(
195        Rc<dyn PluginLifecycle>,
196        ManagedTaskScope,
197        ManagedResourceScope,
198    )> {
199        self.generation.borrow().as_ref().map(|generation| {
200            (
201                generation.lifecycle.clone(),
202                generation.tasks.clone(),
203                generation.resources.clone(),
204            )
205        })
206    }
207}
208
209#[derive(Clone, Debug)]
210pub(super) struct PluginSupervision {
211    pub(super) policy: RestartPolicy,
212    pub(super) criticality: PluginCriticality,
213    pub(super) required_path: bool,
214    pub(super) generation: u64,
215    pub(super) attempts: Vec<Duration>,
216    pub(super) stable_since: Option<Duration>,
217    pub(super) restarting: bool,
218}
219
220#[derive(Debug, Default)]
221pub(super) struct ShutdownCoordinator {
222    pub(super) started: Cell<bool>,
223    pub(super) cleanup_started_at: Cell<Option<Duration>>,
224    pub(super) completed: Cell<bool>,
225    pub(super) outcome: RefCell<Option<ShutdownOutcome>>,
226    pub(super) waiters: RefCell<Vec<oneshot::Sender<ShutdownOutcome>>>,
227}
228
229impl ShutdownCoordinator {
230    pub(super) fn start(&self, started_at: Duration) -> bool {
231        if self.started.replace(true) {
232            return false;
233        }
234        self.cleanup_started_at.set(Some(started_at));
235        true
236    }
237
238    pub(super) fn begin_completion(&self) -> bool {
239        !self.completed.replace(true)
240    }
241
242    pub(super) fn publish(&self, outcome: &ShutdownOutcome) {
243        self.outcome.replace(Some(outcome.clone()));
244        for waiter in self.waiters.borrow_mut().drain(..) {
245            let _ = waiter.send(outcome.clone());
246        }
247    }
248
249    pub(super) fn wait(&self) -> LocalBoxFuture<'static, ShutdownOutcome> {
250        if let Some(outcome) = self.outcome.borrow().clone() {
251            return Box::pin(futures::future::ready(outcome));
252        }
253        let (complete, waiter) = oneshot::channel();
254        self.waiters.borrow_mut().push(complete);
255        Box::pin(async move {
256            waiter.await.unwrap_or(ShutdownOutcome::RuntimeFailure {
257                error: RuntimeFailure::Internal {
258                    detail: "shutdown coordinator terminated before publishing an outcome"
259                        .to_owned(),
260                },
261            })
262        })
263    }
264}
265
266pub(super) struct NativeAppRuntime {
267    pub(super) startup_context: RefCell<Option<InvocationContext>>,
268    pub(super) startup_cleanup: Option<super::cleanup::StartupCleanupBudget>,
269    pub(super) cleanup_timeout: Option<Duration>,
270    pub(super) executions: Rc<super::settlement::ExecutionLedger>,
271    pub(super) plan: ResolvedAppPlan,
272    pub(super) adapters: Rc<ExecutionAdapterCatalog>,
273    pub(super) plugins: BTreeMap<String, NativePluginRuntime>,
274    pub(super) dependencies: BTreeMap<String, PluginDependencies>,
275    pub(super) endpoint_states: BTreeMap<(String, String), Rc<NativeEndpointState>>,
276    pub(super) stream_endpoint_states: NativeStreamEndpointStateTable,
277    pub(super) event_endpoint_states: NativeEventEndpointStateTable,
278    pub(super) supervision: RefCell<BTreeMap<String, PluginSupervision>>,
279    pub(super) supervision_tasks: RefCell<BTreeMap<String, ManagedTask>>,
280    pub(super) activation_order: Vec<String>,
281    pub(super) ready_gate: AppReadyGate,
282    pub(super) admission: AppAdmission,
283    pub(super) driver: DriverControl,
284    pub(super) diagnostics: RuntimeDiagnostics,
285    pub(super) request_ids: Rc<Cell<RequestId>>,
286    pub(super) supervision_cancellation: CancellationToken,
287    pub(super) shutdown_started: Cell<bool>,
288    pub(super) shutdown: ShutdownCoordinator,
289    pub(super) shutdown_task: RefCell<Option<DriverTask>>,
290    pub(super) terminal_failure: RefCell<Option<RuntimeFailure>>,
291    pub(super) cleanup_failure: RefCell<Option<RuntimeFailure>>,
292}
293
294impl std::fmt::Debug for NativeAppRuntime {
295    fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
296        formatter
297            .debug_struct("NativeAppRuntime")
298            .field("plugin_count", &self.plugins.len())
299            .field("endpoint_count", &self.endpoint_states.len())
300            .field("stream_endpoint_count", &self.stream_endpoint_states.len())
301            .field("event_endpoint_count", &self.event_endpoint_states.len())
302            .field("ready", &self.ready_gate.is_open())
303            .field("accepting", &self.admission.is_open())
304            .field("next_request_id", &self.request_ids.get())
305            .field("shutdown_started", &self.shutdown_started.get())
306            .field("cleanup_started", &self.shutdown.started.get())
307            .field("cleanup_completed", &self.shutdown.completed.get())
308            .field(
309                "terminal_failure",
310                &self.terminal_failure.borrow().is_some(),
311            )
312            .finish_non_exhaustive()
313    }
314}
315
316impl NativeAppRuntime {
317    pub(super) fn record_cleanup_failure(&self, error: &RuntimeFailure) {
318        // Optional Plugin failure need not close admission, but retirement must
319        // retain cleanup evidence even after that generation has been removed.
320        self.cleanup_failure
321            .borrow_mut()
322            .get_or_insert_with(|| error.clone());
323    }
324
325    pub(super) fn mark_plugin_endpoints_unavailable(&self, instance_key: &str) {
326        for ((provider, _), endpoint) in &self.endpoint_states {
327            if provider == instance_key {
328                endpoint.mark_unavailable();
329            }
330        }
331        for ((provider, _), endpoint) in &self.stream_endpoint_states {
332            if provider == instance_key {
333                endpoint.mark_unavailable();
334            }
335        }
336        for ((provider, _), endpoint) in &self.event_endpoint_states {
337            if provider == instance_key {
338                endpoint.mark_unavailable();
339            }
340        }
341    }
342
343    pub(super) fn begin_shutdown(&self) {
344        let admission_closed_at = (self.driver.now)();
345        if self.shutdown_started.replace(true) {
346            return;
347        }
348        self.admission.close();
349        self.supervision_cancellation.cancel();
350        for endpoint in self.endpoint_states.values() {
351            endpoint.cancel();
352        }
353        for endpoint in self.stream_endpoint_states.values() {
354            endpoint.cancel();
355        }
356        for endpoint in self.event_endpoint_states.values() {
357            endpoint.cancel();
358        }
359        for plugin in self.plugins.values() {
360            if let Some((_, tasks, resources)) = plugin.generation_parts() {
361                tasks.close();
362                resources.close();
363            }
364        }
365        self.diagnostics
366            .emit(DiagnosticSource::Shutdown, admission_closed_at, |_| {
367                DiagnosticEvent::ShutdownAdmissionClosed
368            });
369    }
370
371    pub(super) fn complete_shutdown(&self, outcome: &ShutdownOutcome) {
372        if !self.shutdown.begin_completion() {
373            return;
374        }
375        let completed_at = (self.driver.now)();
376        let started_at = self
377            .shutdown
378            .cleanup_started_at
379            .get()
380            .unwrap_or(completed_at);
381        let diagnostic_outcome = match outcome {
382            ShutdownOutcome::Clean => DiagnosticShutdownOutcome::Clean,
383            ShutdownOutcome::RuntimeFailure { .. } => DiagnosticShutdownOutcome::RuntimeFailure,
384            ShutdownOutcome::Timeout => DiagnosticShutdownOutcome::Timeout,
385        };
386        self.diagnostics
387            .emit(DiagnosticSource::Shutdown, completed_at, |_| {
388                DiagnosticEvent::ShutdownCompleted {
389                    outcome: diagnostic_outcome,
390                    elapsed: completed_at.saturating_sub(started_at),
391                }
392            });
393        if let ShutdownOutcome::RuntimeFailure { error } = outcome {
394            self.diagnostics
395                .emit_runtime_failure(completed_at, None, error);
396        }
397        self.shutdown.publish(outcome);
398    }
399}
400
401/// A started native App whose generated clients can invoke resolved bindings.
402#[derive(Clone, Debug)]
403pub struct NativeApp {
404    pub(super) bindings: BTreeMap<(String, &'static str), Vec<NativeEndpointBinding>>,
405    pub(super) stream_bindings: NativeStreamBindingTable,
406    pub(super) event_bindings: NativeEventBindingTable,
407    pub(super) diagnostics: RuntimeDiagnostics,
408    pub(super) runtime: Rc<NativeAppRuntime>,
409}
410
411impl NativeApp {
412    fn diagnostic_failure<T>(
413        &self,
414        instance_key: Option<&str>,
415        error: RuntimeFailure,
416    ) -> Result<T, RuntimeFailure> {
417        let instance_key = instance_key
418            .filter(|instance_key| self.runtime.plan.plugin_instance(instance_key).is_some());
419        self.runtime.diagnostics.emit_runtime_failure(
420            (self.runtime.driver.now)(),
421            instance_key,
422            &error,
423        );
424        Err(error)
425    }
426
427    /// Confirms that a generated client has one resolved binding before use.
428    pub fn ensure_binding<C: RequestCapability>(
429        &self,
430        caller_instance: &str,
431    ) -> Result<(), RuntimeFailure> {
432        if self.runtime.admission.is_closed() {
433            return self.diagnostic_failure(Some(caller_instance), RuntimeFailure::AdmissionClosed);
434        }
435        if self
436            .endpoints::<C>(caller_instance)
437            .is_some_and(|endpoints| !endpoints.is_empty())
438        {
439            return Ok(());
440        }
441        self.diagnostic_failure(
442            Some(caller_instance),
443            RuntimeFailure::Unavailable { capability: C::ID },
444        )
445    }
446
447    /// Materializes one typed handle from the immutable binding selected by the Plan.
448    pub fn handle<C: RequestCapability>(
449        &self,
450        caller_instance: &str,
451    ) -> Result<NativeRequestHandle<C>, RuntimeFailure> {
452        self.validate_requirement_lookup(caller_instance, C::ID)?;
453        if self.runtime.admission.is_closed() {
454            return self.diagnostic_failure(Some(caller_instance), RuntimeFailure::AdmissionClosed);
455        }
456        let Some(endpoints) = self
457            .endpoints::<C>(caller_instance)
458            .filter(|endpoints| !endpoints.is_empty())
459        else {
460            return self.diagnostic_failure(
461                Some(caller_instance),
462                RuntimeFailure::Unavailable { capability: C::ID },
463            );
464        };
465        Ok(NativeRequestHandle::from_endpoints(
466            endpoints,
467            self.runtime.clone(),
468            caller_instance,
469            false,
470        ))
471    }
472
473    /// Materializes an optional typed handle; an absent binding remains `None`.
474    pub fn optional_handle<C: RequestCapability>(
475        &self,
476        caller_instance: &str,
477    ) -> Option<NativeRequestHandle<C>> {
478        self.validate_requirement_lookup(caller_instance, C::ID)
479            .ok()?;
480        let caller_instance = caller_instance.to_owned();
481        self.endpoints::<C>(&caller_instance)
482            .filter(|endpoints| !endpoints.is_empty())
483            .map(|endpoints| {
484                NativeRequestHandle::from_endpoints(
485                    endpoints,
486                    self.runtime.clone(),
487                    &caller_instance,
488                    false,
489                )
490            })
491    }
492
493    /// Materializes a typed handle whose endpoints may be empty for a `many` requirement.
494    pub fn many_handle<C: RequestCapability>(
495        &self,
496        caller_instance: &str,
497    ) -> Result<NativeRequestHandle<C>, RuntimeFailure> {
498        self.validate_requirement_lookup(caller_instance, C::ID)?;
499        if self.runtime.admission.is_closed() {
500            return self.diagnostic_failure(Some(caller_instance), RuntimeFailure::AdmissionClosed);
501        }
502        let endpoints = self.endpoints::<C>(caller_instance).unwrap_or(&[]);
503        Ok(NativeRequestHandle::from_endpoints(
504            endpoints,
505            self.runtime.clone(),
506            caller_instance,
507            false,
508        ))
509    }
510
511    /// Returns the number of immutable provider endpoints bound to one requirement.
512    pub fn binding_count<C: RequestCapability>(&self, caller_instance: &str) -> usize {
513        self.endpoints::<C>(caller_instance).map_or(0, <[_]>::len)
514    }
515
516    /// Returns whether every declared Plugin has completed activation.
517    pub fn is_ready(&self) -> bool {
518        self.runtime.ready_gate.is_open()
519    }
520
521    /// Returns the App-wide readiness signal observed by Plugin tasks.
522    pub fn ready_gate(&self) -> AppReadyGate {
523        self.runtime.ready_gate.clone()
524    }
525
526    /// Returns whether new externally triggered work may be admitted.
527    pub fn is_accepting(&self) -> bool {
528        self.runtime.admission.is_open()
529    }
530
531    /// Returns the App-wide admission state.
532    pub fn admission(&self) -> AppAdmission {
533        self.runtime.admission.clone()
534    }
535
536    /// Returns the opt-in Runtime Diagnostics port for this App.
537    pub fn diagnostics(&self) -> RuntimeDiagnostics {
538        self.diagnostics.clone()
539    }
540
541    /// Returns the exact resolved Capability dependencies for one Plugin Instance.
542    ///
543    /// Runners use this host-facing view when they must preserve provider identity
544    /// while transferring one binding across an Execution Lane. Ordinary Plugin
545    /// code receives the same Interface through its lifecycle context.
546    pub fn dependencies(
547        &self,
548        caller_instance: &str,
549    ) -> Result<PluginDependencies, RuntimeFailure> {
550        if self.runtime.admission.is_closed() {
551            return self.diagnostic_failure(Some(caller_instance), RuntimeFailure::AdmissionClosed);
552        }
553        self.runtime
554            .dependencies
555            .get(caller_instance)
556            .cloned()
557            .ok_or_else(|| RuntimeFailure::InvalidResolvedPlan {
558                detail: format!(
559                    "Plugin Instance `{caller_instance}` has no resolved dependency table"
560                ),
561            })
562    }
563
564    /// Returns current bounded request queue depths grouped by provider Instance.
565    ///
566    /// This structural snapshot contains no request payloads and is intended for
567    /// Runner-owned placement diagnostics.
568    pub fn instance_queue_depths(&self) -> BTreeMap<String, usize> {
569        let mut depths = BTreeMap::new();
570        for endpoints in self.bindings.values() {
571            for endpoint in endpoints {
572                let depth = endpoint
573                    .admissions
574                    .values()
575                    .map(RequestAdmission::queue_depth)
576                    .sum::<usize>();
577                *depths.entry(endpoint.plugin_instance.clone()).or_insert(0) += depth;
578            }
579        }
580        depths
581    }
582
583    /// Returns the terminal supervision failure, when a critical App path exhausted its budget.
584    pub fn terminal_failure(&self) -> Option<RuntimeFailure> {
585        self.runtime.terminal_failure.borrow().clone()
586    }
587
588    /// Returns whether supervision has produced a terminal App failure.
589    pub fn is_failed(&self) -> bool {
590        self.runtime.terminal_failure.borrow().is_some()
591    }
592
593    /// Returns the current ready generation for one Plugin Instance, when it is available.
594    pub fn plugin_generation(&self, instance_key: &str) -> Option<u64> {
595        self.runtime
596            .supervision
597            .borrow()
598            .get(instance_key)
599            .and_then(|state| {
600                let request_current =
601                    self.runtime
602                        .endpoint_states
603                        .iter()
604                        .any(|((plugin, _), endpoint)| {
605                            plugin == instance_key && endpoint.is_current(state.generation)
606                        });
607                let stream_current =
608                    self.runtime
609                        .stream_endpoint_states
610                        .iter()
611                        .any(|((plugin, _), endpoint)| {
612                            plugin == instance_key && endpoint.is_current(state.generation)
613                        });
614                let event_current =
615                    self.runtime
616                        .event_endpoint_states
617                        .iter()
618                        .any(|((plugin, _), endpoint)| {
619                            plugin == instance_key && endpoint.is_current(state.generation)
620                        });
621                (request_current || stream_current || event_current).then_some(state.generation)
622            })
623    }
624
625    /// Reports a Plugin Instance failure and schedules its finite supervision policy.
626    pub fn report_plugin_failure(&self, instance_key: &str) -> Result<(), RuntimeFailure> {
627        if !begin_plugin_supervision(&self.runtime, instance_key)? {
628            return Ok(());
629        }
630        schedule_plugin_supervision(&self.runtime, instance_key).map_err(|error| {
631            handle_supervision_schedule_failure(&self.runtime, instance_key, error)
632        })
633    }
634
635    /// Starts shutdown admission closure and cooperative cancellation.
636    pub fn request_shutdown(&self) {
637        self.runtime.begin_shutdown();
638    }
639
640    /// Performs bounded graceful shutdown using one global deadline.
641    pub async fn shutdown(&self, timeout: Duration) -> ShutdownOutcome {
642        self.shutdown_with_budget(super::cleanup::CleanupBudget::after(
643            &self.runtime.driver,
644            timeout,
645        ))
646        .await
647    }
648
649    pub(super) async fn shutdown_with_budget(
650        &self,
651        budget: super::cleanup::CleanupBudget,
652    ) -> ShutdownOutcome {
653        self.runtime.begin_shutdown();
654        let cleanup_started_at = (self.runtime.driver.now)();
655        if self.runtime.shutdown.start(cleanup_started_at) {
656            let timeout = budget.remaining();
657            self.runtime
658                .diagnostics
659                .emit(DiagnosticSource::Shutdown, cleanup_started_at, |_| {
660                    DiagnosticEvent::ShutdownCleanupStarted { timeout }
661                });
662            let runtime = self.runtime.clone();
663            let worker_runtime = runtime.clone();
664            match (runtime.driver.spawn_local)(Box::pin(async move {
665                let outcome = shutdown_native_plugins(&worker_runtime, budget).await;
666                worker_runtime.complete_shutdown(&outcome);
667            })) {
668                Ok(task) => {
669                    runtime.shutdown_task.replace(Some(task));
670                }
671                Err(error) => {
672                    runtime.complete_shutdown(&ShutdownOutcome::RuntimeFailure {
673                        error: RuntimeFailure::Internal {
674                            detail: format!("failed to schedule App shutdown: {error:?}"),
675                        },
676                    });
677                }
678            }
679        }
680        self.runtime.shutdown.wait().await
681    }
682
683    /// Invokes a generated request Operation through the caller's resolved binding.
684    pub async fn invoke<C: RequestCapability>(
685        &self,
686        caller_instance: &str,
687        operation: &str,
688        request: C::Request,
689    ) -> Result<Result<C::Response, C::DomainError>, RuntimeFailure> {
690        self.handle::<C>(caller_instance)?
691            .invoke(operation, request)
692            .await
693    }
694
695    /// Creates a request context with a fresh Kernel Request ID.
696    ///
697    /// `deadline` is an absolute instant returned by the selected
698    /// [`RuntimeDriver`]'s monotonic clock.
699    pub fn invocation_context(
700        &self,
701        deadline: Option<Duration>,
702        cancellation: CancellationToken,
703    ) -> InvocationContext {
704        InvocationContext::new(self.next_request_id(), deadline, cancellation)
705    }
706
707    /// Creates a request context whose deadline is relative to the Driver's clock.
708    pub fn invocation_context_after(
709        &self,
710        timeout: Duration,
711        cancellation: CancellationToken,
712    ) -> InvocationContext {
713        self.invocation_context(
714            Some((self.runtime.driver.now)().saturating_add(timeout)),
715            cancellation,
716        )
717    }
718
719    /// Invokes a request with an explicit propagated Invocation Context.
720    pub async fn invoke_with_context<C: RequestCapability>(
721        &self,
722        caller_instance: &str,
723        operation: &str,
724        context: InvocationContext,
725        request: C::Request,
726    ) -> Result<Result<C::Response, C::DomainError>, RuntimeFailure> {
727        self.handle::<C>(caller_instance)?
728            .invoke_with_context(operation, context, request)
729            .await
730    }
731
732    /// Materializes one typed bidirectional stream handle from the resolved Plan.
733    pub fn stream_handle<C: StreamCapability>(
734        &self,
735        caller_instance: &str,
736    ) -> Result<NativeStreamHandle<C>, RuntimeFailure> {
737        self.validate_requirement_lookup(caller_instance, C::ID)?;
738        if self.runtime.admission.is_closed() {
739            return self.diagnostic_failure(Some(caller_instance), RuntimeFailure::AdmissionClosed);
740        }
741        let Some(endpoints) = self
742            .stream_endpoints::<C>(caller_instance)
743            .filter(|endpoints| !endpoints.is_empty())
744        else {
745            return self.diagnostic_failure(
746                Some(caller_instance),
747                RuntimeFailure::Unavailable { capability: C::ID },
748            );
749        };
750        Ok(NativeStreamHandle::from_endpoints(
751            endpoints,
752            self.runtime.clone(),
753            caller_instance,
754            false,
755        ))
756    }
757
758    /// Materializes an optional typed bidirectional stream handle.
759    pub fn optional_stream_handle<C: StreamCapability>(
760        &self,
761        caller_instance: &str,
762    ) -> Option<NativeStreamHandle<C>> {
763        self.validate_requirement_lookup(caller_instance, C::ID)
764            .ok()?;
765        let caller_instance = caller_instance.to_owned();
766        self.stream_endpoints::<C>(&caller_instance)
767            .filter(|endpoints| !endpoints.is_empty())
768            .map(|endpoints| {
769                NativeStreamHandle::from_endpoints(
770                    endpoints,
771                    self.runtime.clone(),
772                    &caller_instance,
773                    false,
774                )
775            })
776    }
777
778    /// Returns the number of immutable stream endpoints bound to one requirement.
779    pub fn stream_binding_count<C: StreamCapability>(&self, caller_instance: &str) -> usize {
780        self.stream_endpoints::<C>(caller_instance)
781            .map_or(0, <[_]>::len)
782    }
783
784    /// Materializes a typed Event handle and requires at least one subscriber.
785    pub fn event_handle<C: EventCapability>(
786        &self,
787        caller_instance: &str,
788    ) -> Result<NativeEventHandle<C>, RuntimeFailure> {
789        self.validate_requirement_lookup(caller_instance, C::ID)?;
790        if self.runtime.admission.is_closed() {
791            return self.diagnostic_failure(Some(caller_instance), RuntimeFailure::AdmissionClosed);
792        }
793        let Some(endpoints) = self
794            .event_endpoints::<C>(caller_instance)
795            .filter(|endpoints| !endpoints.is_empty())
796        else {
797            return self.diagnostic_failure(
798                Some(caller_instance),
799                RuntimeFailure::Unavailable { capability: C::ID },
800            );
801        };
802        Ok(NativeEventHandle::from_endpoints(
803            endpoints,
804            self.runtime.clone(),
805            caller_instance,
806            false,
807        ))
808    }
809
810    /// Materializes an optional typed Event handle.
811    pub fn optional_event_handle<C: EventCapability>(
812        &self,
813        caller_instance: &str,
814    ) -> Option<NativeEventHandle<C>> {
815        self.validate_requirement_lookup(caller_instance, C::ID)
816            .ok()?;
817        let caller_instance = caller_instance.to_owned();
818        self.event_endpoints::<C>(&caller_instance)
819            .filter(|endpoints| !endpoints.is_empty())
820            .map(|endpoints| {
821                NativeEventHandle::from_endpoints(
822                    endpoints,
823                    self.runtime.clone(),
824                    &caller_instance,
825                    false,
826                )
827            })
828    }
829
830    /// Materializes a typed Event handle whose endpoint set may be empty.
831    pub fn many_event_handle<C: EventCapability>(
832        &self,
833        caller_instance: &str,
834    ) -> Result<NativeEventHandle<C>, RuntimeFailure> {
835        self.validate_requirement_lookup(caller_instance, C::ID)?;
836        if self.runtime.admission.is_closed() {
837            return self.diagnostic_failure(Some(caller_instance), RuntimeFailure::AdmissionClosed);
838        }
839        let endpoints = self.event_endpoints::<C>(caller_instance).unwrap_or(&[]);
840        Ok(NativeEventHandle::from_endpoints(
841            endpoints,
842            self.runtime.clone(),
843            caller_instance,
844            false,
845        ))
846    }
847
848    /// Returns the number of immutable Event subscriber endpoints bound to a requirement.
849    pub fn event_binding_count<C: EventCapability>(&self, caller_instance: &str) -> usize {
850        self.event_endpoints::<C>(caller_instance)
851            .map_or(0, <[_]>::len)
852    }
853
854    fn validate_requirement_lookup(
855        &self,
856        caller: &str,
857        capability: &'static str,
858    ) -> Result<(), RuntimeFailure> {
859        let declarations = self
860            .runtime
861            .plan
862            .plugin_instance(caller)
863            .map_or(0, |instance| {
864                instance
865                    .required_capabilities()
866                    .iter()
867                    .filter(|requirement| requirement.capability_id() == capability)
868                    .count()
869            });
870        if declarations > 1 {
871            return Err(RuntimeFailure::AmbiguousBinding {
872                capability,
873                providers: declarations,
874            });
875        }
876        Ok(())
877    }
878
879    pub(super) fn next_request_id(&self) -> RequestId {
880        let request_id = self.runtime.request_ids.get();
881        self.runtime.request_ids.set(request_id.saturating_add(1));
882        request_id
883    }
884
885    pub(super) fn endpoints<C: RequestCapability>(
886        &self,
887        caller_instance: &str,
888    ) -> Option<&[NativeEndpointBinding]> {
889        self.bindings
890            .get(&(caller_instance.to_owned(), C::ID))
891            .map(Vec::as_slice)
892    }
893
894    pub(super) fn stream_endpoints<C: StreamCapability>(
895        &self,
896        caller_instance: &str,
897    ) -> Option<&[NativeStreamEndpointBinding]> {
898        self.stream_bindings
899            .get(&(caller_instance.to_owned(), C::ID))
900            .map(Vec::as_slice)
901    }
902
903    pub(super) fn event_endpoints<C: EventCapability>(
904        &self,
905        caller_instance: &str,
906    ) -> Option<&[event::NativeEventEndpointBinding]> {
907        self.event_bindings
908            .get(&(caller_instance.to_owned(), C::ID))
909            .map(Vec::as_slice)
910    }
911}