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 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 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#[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 pub fn plan_snapshot(&self) -> ResolvedAppPlan {
428 self.runtime.snapshot.borrow().as_ref().clone()
429 }
430
431 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 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 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 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 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 pub fn binding_count<C: RequestCapability>(&self, caller_instance: &str) -> usize {
536 self.endpoints::<C>(caller_instance).map_or(0, <[_]>::len)
537 }
538
539 pub fn is_ready(&self) -> bool {
541 self.runtime.ready_gate.is_open()
542 }
543
544 pub fn ready_gate(&self) -> AppReadyGate {
546 self.runtime.ready_gate.clone()
547 }
548
549 pub fn is_accepting(&self) -> bool {
551 self.runtime.admission.is_open()
552 }
553
554 pub fn admission(&self) -> AppAdmission {
556 self.runtime.admission.clone()
557 }
558
559 pub fn diagnostics(&self) -> RuntimeDiagnostics {
561 self.diagnostics.clone()
562 }
563
564 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 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 pub fn terminal_failure(&self) -> Option<RuntimeFailure> {
608 self.runtime.terminal_failure.borrow().clone()
609 }
610
611 pub fn is_failed(&self) -> bool {
613 self.runtime.terminal_failure.borrow().is_some()
614 }
615
616 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 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 pub fn request_shutdown(&self) {
660 self.runtime.begin_shutdown();
661 }
662
663 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 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 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 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 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 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 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 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 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 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 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 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}