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