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 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#[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 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 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 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 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 pub fn binding_count<C: RequestCapability>(&self, caller_instance: &str) -> usize {
513 self.endpoints::<C>(caller_instance).map_or(0, <[_]>::len)
514 }
515
516 pub fn is_ready(&self) -> bool {
518 self.runtime.ready_gate.is_open()
519 }
520
521 pub fn ready_gate(&self) -> AppReadyGate {
523 self.runtime.ready_gate.clone()
524 }
525
526 pub fn is_accepting(&self) -> bool {
528 self.runtime.admission.is_open()
529 }
530
531 pub fn admission(&self) -> AppAdmission {
533 self.runtime.admission.clone()
534 }
535
536 pub fn diagnostics(&self) -> RuntimeDiagnostics {
538 self.diagnostics.clone()
539 }
540
541 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 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 pub fn terminal_failure(&self) -> Option<RuntimeFailure> {
585 self.runtime.terminal_failure.borrow().clone()
586 }
587
588 pub fn is_failed(&self) -> bool {
590 self.runtime.terminal_failure.borrow().is_some()
591 }
592
593 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 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 pub fn request_shutdown(&self) {
637 self.runtime.begin_shutdown();
638 }
639
640 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 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 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 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 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 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 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 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 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 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 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 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}