1use super::{
2 ActivateContext, AppAdmission, AppReadyGate, BTreeMap, CancellationToken, Cell,
3 DeactivationReason, DriverControl, ExecutionAdapterCatalog, ManagedResourceScope,
4 ManagedTaskScope, NativeApp, NativeAppRuntime, NativeBindingTable, NativeEndpointBinding,
5 NativeEndpointState, NativeEndpointStateTable, NativeEventBindingTable,
6 NativeEventEndpointStateTable, NativeExecutionAdapter, NativePluginGeneration,
7 NativePluginRuntime, NativeStreamBindingTable, NativeStreamEndpointBinding,
8 NativeStreamEndpointState, NativeStreamEndpointStateTable, PlanResolutionError,
9 PluginDependencies, PluginDependency, PluginDependencyHandle, PluginEventDependencyHandle,
10 PluginStreamDependencyHandle, PrepareContext, PreparedBinding, PreparedEventBinding,
11 PreparedNativeApp, PreparedNativePlugin, PreparedStreamBinding, Rc, RefCell, RequestAdmission,
12 ResolvedAppPlan, RuntimeDiagnostics, RuntimeDriver, RuntimeFailure, ShutdownCoordinator, Weak,
13 begin_plugin_supervision, deactivate_in_reverse, event, handle_supervision_schedule_failure,
14 plugin_supervision, schedule_plugin_supervision, validate_native_endpoint_set,
15};
16
17#[derive(Clone, Debug, Eq, PartialEq)]
19pub enum PlanValidationError {
20 UnsupportedSchemaVersion { expected: u32, actual: u32 },
22 InvalidResolvedPlan { detail: String },
24}
25
26impl std::fmt::Display for PlanValidationError {
27 fn fmt(&self, formatter: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
28 match self {
29 Self::UnsupportedSchemaVersion { expected, actual } => write!(
30 formatter,
31 "unsupported Resolved App Plan schema {actual}; expected {expected}"
32 ),
33 Self::InvalidResolvedPlan { detail } => {
34 write!(formatter, "invalid Resolved App Plan: {detail}")
35 }
36 }
37 }
38}
39
40impl std::error::Error for PlanValidationError {}
41
42#[derive(Debug)]
44pub struct Kernel;
45
46impl Kernel {
47 pub async fn start_native<D: RuntimeDriver, A: NativeExecutionAdapter>(
49 plan: ResolvedAppPlan,
50 driver: D,
51 adapter: A,
52 ) -> Result<NativeApp, RuntimeFailure> {
53 Self::start_native_with_diagnostics(plan, driver, adapter, RuntimeDiagnostics::new()).await
54 }
55
56 pub async fn start_native_with_diagnostics<D: RuntimeDriver, A: NativeExecutionAdapter>(
58 plan: ResolvedAppPlan,
59 driver: D,
60 adapter: A,
61 diagnostics: RuntimeDiagnostics,
62 ) -> Result<NativeApp, RuntimeFailure> {
63 Self::start_with_diagnostics(
64 plan,
65 driver,
66 ExecutionAdapterCatalog::single(adapter),
67 diagnostics,
68 )
69 .await
70 }
71
72 pub async fn start<D: RuntimeDriver>(
74 plan: ResolvedAppPlan,
75 driver: D,
76 adapters: ExecutionAdapterCatalog,
77 ) -> Result<NativeApp, RuntimeFailure> {
78 Self::start_with_diagnostics(plan, driver, adapters, RuntimeDiagnostics::new()).await
79 }
80
81 #[allow(
83 clippy::too_many_lines,
84 reason = "startup remains linear so validation, preparation, and activation fail closed in order"
85 )]
86 pub async fn start_with_diagnostics<D: RuntimeDriver>(
87 plan: ResolvedAppPlan,
88 driver: D,
89 adapters: ExecutionAdapterCatalog,
90 diagnostics: RuntimeDiagnostics,
91 ) -> Result<NativeApp, RuntimeFailure> {
92 if plan
93 .plugin_instances()
94 .iter()
95 .any(|instance| instance.authoring_version() == 2)
96 {
97 return Self::start_controlled(
98 plan,
99 driver,
100 adapters,
101 diagnostics,
102 super::InvocationContext::new(0, None, CancellationToken::new()),
103 super::startup::DEFAULT_STARTUP_CLEANUP_TIMEOUT,
104 )
105 .await;
106 }
107 Self::start_owned(plan, driver, adapters, diagnostics, None, None).await
108 }
109
110 pub async fn start_controlled<D: RuntimeDriver>(
113 plan: ResolvedAppPlan,
114 driver: D,
115 adapters: ExecutionAdapterCatalog,
116 diagnostics: RuntimeDiagnostics,
117 context: super::InvocationContext,
118 cleanup_timeout: std::time::Duration,
119 ) -> Result<NativeApp, RuntimeFailure> {
120 super::startup::start(
121 plan,
122 driver,
123 adapters,
124 diagnostics,
125 context,
126 cleanup_timeout,
127 )
128 .await
129 }
130
131 #[allow(
132 clippy::too_many_lines,
133 reason = "startup remains linear so validation, preparation and activation fail closed"
134 )]
135 pub(super) async fn start_owned<D: RuntimeDriver>(
136 plan: ResolvedAppPlan,
137 driver: D,
138 adapters: ExecutionAdapterCatalog,
139 diagnostics: RuntimeDiagnostics,
140 startup_context: Option<super::InvocationContext>,
141 startup_cleanup: Option<super::cleanup::StartupCleanupBudget>,
142 ) -> Result<NativeApp, RuntimeFailure> {
143 let activation_order = match plan.activation_order() {
146 Ok(order) => order,
147 Err(error) => {
148 let error = runtime_plan_error(&error);
149 diagnostics.emit_runtime_failure(driver.now(), None, &error);
150 return Err(error);
151 }
152 };
153 let adapters = Rc::new(adapters);
154 let PreparedNativeApp {
155 bindings: prepared_bindings,
156 stream_bindings: prepared_stream_bindings,
157 event_bindings: prepared_event_bindings,
158 generations,
159 } = match adapters.prepare(&plan) {
160 Ok(prepared) => prepared,
161 Err(error) => {
162 diagnostics.emit_runtime_failure(driver.now(), None, &error);
163 return Err(error);
164 }
165 };
166 if let Err(error) = validate_prepared_native_app(
167 &plan,
168 &prepared_bindings,
169 &prepared_stream_bindings,
170 &prepared_event_bindings,
171 &generations,
172 ) {
173 diagnostics.emit_runtime_failure(driver.now(), None, &error);
174 return Err(error);
175 }
176 let (bindings, endpoint_states) = native_bindings(&plan, &prepared_bindings);
177 let (stream_bindings, stream_endpoint_states) =
178 native_stream_bindings(&plan, &prepared_stream_bindings);
179 let (event_bindings, event_endpoint_states) =
180 native_event_bindings(&plan, &prepared_event_bindings);
181 let runtime_link = Rc::new(RefCell::new(Weak::new()));
182 let dependencies = plugin_dependencies(
183 &plan,
184 &bindings,
185 &stream_bindings,
186 &event_bindings,
187 &runtime_link,
188 );
189 let driver_control = DriverControl::new(&driver);
190 let admission = AppAdmission::new();
191 let startup_cancellation = startup_context
192 .as_ref()
193 .map(super::InvocationContext::cancellation);
194 let plugin_runtimes =
195 native_plugin_runtimes(&plan, &driver, generations, startup_cancellation.as_ref());
196 let ready_gate = AppReadyGate::new();
197 let supervision = plugin_supervision(&plan);
198 let cleanup_timeout = startup_cleanup
199 .as_ref()
200 .map(super::cleanup::StartupCleanupBudget::timeout);
201 let runtime = Rc::new(NativeAppRuntime {
202 startup_context: RefCell::new(startup_context),
203 startup_cleanup,
204 cleanup_timeout,
205 executions: Rc::default(),
206 plan,
207 adapters,
208 plugins: plugin_runtimes,
209 dependencies,
210 endpoint_states,
211 stream_endpoint_states,
212 event_endpoint_states,
213 supervision: RefCell::new(supervision),
214 supervision_tasks: RefCell::new(BTreeMap::new()),
215 activation_order,
216 ready_gate,
217 admission,
218 driver: driver_control,
219 diagnostics: diagnostics.clone(),
220 request_ids: Rc::new(Cell::new(1)),
221 supervision_cancellation: CancellationToken::new(),
222 shutdown_started: Cell::new(false),
223 shutdown: ShutdownCoordinator::default(),
224 shutdown_task: RefCell::new(None),
225 terminal_failure: RefCell::new(None),
226 });
227 runtime_link.replace(Rc::downgrade(&runtime));
228 attach_managed_task_failure_handlers(&runtime);
229 runtime.diagnostics.emit(
230 super::DiagnosticSource::Lifecycle,
231 (runtime.driver.now)(),
232 |_| super::DiagnosticEvent::AppStarted {
233 plugin_count: runtime.plan.plugin_instances().len(),
234 },
235 );
236 let prepared_instances = prepare_native_plugins(&runtime).await?;
237 if let Err(error) = construct_native_plugins(&runtime).await {
238 let cleanup_error = deactivate_in_reverse(
239 &runtime.plugins,
240 &runtime.dependencies,
241 &prepared_instances,
242 DeactivationReason::StartupRollback,
243 &runtime.admission,
244 &runtime.diagnostics,
245 &runtime.driver,
246 runtime
247 .startup_cleanup
248 .as_ref()
249 .map(super::cleanup::StartupCleanupBudget::establish),
250 )
251 .await;
252 retain_unsafe_startup(&runtime, cleanup_error.as_ref());
253 runtime
254 .diagnostics
255 .emit_runtime_failure((runtime.driver.now)(), None, &error);
256 return Err(error);
257 }
258 if let Err(error) = activate_native_plugins(&runtime).await {
259 let cleanup_error = deactivate_in_reverse(
260 &runtime.plugins,
261 &runtime.dependencies,
262 &prepared_instances,
263 DeactivationReason::StartupRollback,
264 &runtime.admission,
265 &runtime.diagnostics,
266 &runtime.driver,
267 runtime
268 .startup_cleanup
269 .as_ref()
270 .map(super::cleanup::StartupCleanupBudget::establish),
271 )
272 .await;
273 retain_unsafe_startup(&runtime, cleanup_error.as_ref());
274 runtime
275 .diagnostics
276 .emit_runtime_failure((runtime.driver.now)(), None, &error);
277 return Err(error);
278 }
279 open_native_readiness(&runtime).await;
280 Ok(NativeApp {
281 bindings,
282 stream_bindings,
283 event_bindings,
284 diagnostics,
285 runtime,
286 })
287 }
288}
289
290pub(super) fn attach_managed_task_failure_handlers(runtime: &Rc<NativeAppRuntime>) {
291 for (instance_key, plugin) in &runtime.plugins {
292 let Some((_, tasks, _)) = plugin.generation_parts() else {
293 continue;
294 };
295 attach_managed_task_failure_handler(runtime, instance_key, &tasks);
296 }
297}
298
299fn startup_active(runtime: &NativeAppRuntime) -> Result<(), RuntimeFailure> {
300 let result = runtime
301 .startup_context
302 .borrow()
303 .as_ref()
304 .map_or(Ok(()), |context| {
305 super::ensure_context_active(&runtime.driver, context)
306 });
307 if result.is_err()
308 && let Some(cleanup) = &runtime.startup_cleanup
309 {
310 let now = (runtime.driver.now)();
311 let cleanup_started_at = runtime
312 .startup_context
313 .borrow()
314 .as_ref()
315 .and_then(super::InvocationContext::deadline)
316 .filter(|deadline| now >= *deadline)
317 .unwrap_or(now);
318 cleanup.establish_at(cleanup_started_at);
319 }
320 result
321}
322
323fn lifecycle_cancellation(
324 _runtime: &NativeAppRuntime,
325 tasks: &ManagedTaskScope,
326) -> CancellationToken {
327 tasks.cancellation()
328}
329
330pub(super) fn attach_managed_task_failure_handler(
331 runtime: &Rc<NativeAppRuntime>,
332 instance_key: &str,
333 tasks: &ManagedTaskScope,
334) {
335 let task_runtime = Rc::downgrade(runtime);
336 let task_instance_key = instance_key.to_owned();
337 let handler: Rc<dyn Fn()> = Rc::new(move || {
338 let Some(runtime) = task_runtime.upgrade() else {
339 return;
340 };
341 if begin_plugin_supervision(&runtime, &task_instance_key).unwrap_or(false)
342 && let Err(error) = schedule_plugin_supervision(&runtime, &task_instance_key)
343 {
344 let _ = handle_supervision_schedule_failure(&runtime, &task_instance_key, error);
345 }
346 });
347 tasks.set_failure_handler(&handler);
348}
349
350pub(super) fn runtime_plan_error(error: &PlanResolutionError) -> RuntimeFailure {
351 RuntimeFailure::InvalidResolvedPlan {
352 detail: error.to_string(),
353 }
354}
355
356#[allow(
357 clippy::too_many_lines,
358 reason = "one fail-closed pass keeps request, stream, event, and generation validation aligned"
359)]
360pub(super) fn validate_prepared_native_app(
361 plan: &ResolvedAppPlan,
362 bindings: &[PreparedBinding],
363 stream_bindings: &[PreparedStreamBinding],
364 event_bindings: &[PreparedEventBinding],
365 generations: &BTreeMap<String, PreparedNativePlugin>,
366) -> Result<(), RuntimeFailure> {
367 if generations.len() != plan.plugin_instances().len() {
368 return Err(RuntimeFailure::InvalidResolvedPlan {
369 detail: format!(
370 "Execution Adapters prepared {} Plugin generations; expected {}",
371 generations.len(),
372 plan.plugin_instances().len()
373 ),
374 });
375 }
376 for instance in plan.plugin_instances() {
377 let generation = generations.get(instance.instance_key()).ok_or_else(|| {
378 RuntimeFailure::InvalidResolvedPlan {
379 detail: format!(
380 "Execution Adapters did not prepare Plugin Instance `{}`",
381 instance.instance_key()
382 ),
383 }
384 })?;
385 validate_native_endpoint_set(
386 instance.instance_key(),
387 instance,
388 generation.endpoints(),
389 generation.stream_endpoints(),
390 generation.event_endpoints(),
391 )?;
392 }
393 if let Some(instance_key) = generations.keys().find(|instance_key| {
394 !plan
395 .plugin_instances()
396 .iter()
397 .any(|instance| instance.instance_key() == instance_key.as_str())
398 }) {
399 return Err(RuntimeFailure::InvalidResolvedPlan {
400 detail: format!("Execution Adapter prepared unknown Plugin Instance `{instance_key}`"),
401 });
402 }
403
404 let expected_request_bindings = plan
405 .capability_bindings()
406 .iter()
407 .filter(|binding| {
408 plan.plugin_instance(binding.provider_instance())
409 .and_then(|provider| {
410 provider
411 .provided_capabilities()
412 .iter()
413 .find(|endpoint| endpoint.capability_id() == binding.capability_id())
414 })
415 .is_some_and(|endpoint| !endpoint.request_operations().is_empty())
416 })
417 .count();
418 let expected_stream_bindings = plan
419 .capability_bindings()
420 .iter()
421 .filter(|binding| {
422 plan.plugin_instance(binding.provider_instance())
423 .and_then(|provider| {
424 provider
425 .provided_capabilities()
426 .iter()
427 .find(|endpoint| endpoint.capability_id() == binding.capability_id())
428 })
429 .is_some_and(|endpoint| !endpoint.stream_operations().is_empty())
430 })
431 .count();
432 let expected_event_bindings = plan
433 .capability_bindings()
434 .iter()
435 .filter(|binding| {
436 plan.plugin_instance(binding.provider_instance())
437 .and_then(|provider| {
438 provider
439 .provided_capabilities()
440 .iter()
441 .find(|endpoint| endpoint.capability_id() == binding.capability_id())
442 })
443 .is_some_and(|endpoint| !endpoint.event_operations().is_empty())
444 })
445 .count();
446 if bindings.len() != expected_request_bindings {
447 return Err(RuntimeFailure::InvalidResolvedPlan {
448 detail: if expected_stream_bindings == 0 && stream_bindings.is_empty() {
449 format!(
450 "Execution Adapters prepared {} bindings; expected {}",
451 bindings.len(),
452 expected_request_bindings
453 )
454 } else {
455 format!(
456 "Execution Adapters prepared {} request bindings; expected {}",
457 bindings.len(),
458 expected_request_bindings
459 )
460 },
461 });
462 }
463 if stream_bindings.len() != expected_stream_bindings {
464 return Err(RuntimeFailure::InvalidResolvedPlan {
465 detail: format!(
466 "Execution Adapters prepared {} stream bindings; expected {}",
467 stream_bindings.len(),
468 expected_stream_bindings
469 ),
470 });
471 }
472 if event_bindings.len() != expected_event_bindings {
473 return Err(RuntimeFailure::InvalidResolvedPlan {
474 detail: format!(
475 "Execution Adapters prepared {} Event bindings; expected {}",
476 event_bindings.len(),
477 expected_event_bindings
478 ),
479 });
480 }
481 for planned in plan.capability_bindings() {
482 let provider = generations
483 .get(planned.provider_instance())
484 .expect("the resolved Plan references one validated provider generation");
485 let descriptor = plan
486 .plugin_instance(planned.provider_instance())
487 .and_then(|provider| {
488 provider
489 .provided_capabilities()
490 .iter()
491 .find(|endpoint| endpoint.capability_id() == planned.capability_id())
492 })
493 .expect("the resolved Plan references one validated provider endpoint");
494 if !descriptor.request_operations().is_empty() {
495 let matching: Vec<_> = bindings
496 .iter()
497 .filter(|prepared| {
498 prepared.requirement_id() == planned.requirement_id()
499 && prepared.consumer_instance == planned.consumer_instance()
500 && prepared.provider_instance == planned.provider_instance()
501 && prepared.endpoint.capability_id() == planned.capability_id()
502 && prepared.endpoint.descriptor_version() == planned.descriptor_version()
503 })
504 .collect();
505 if matching.len() != 1 {
506 return Err(RuntimeFailure::InvalidResolvedPlan {
507 detail: format!(
508 "Execution Adapters prepared {} request bindings for `{}:{}:{}`; expected 1",
509 matching.len(),
510 planned.consumer_instance(),
511 planned.capability_id(),
512 planned.provider_instance()
513 ),
514 });
515 }
516 if !provider
517 .endpoints()
518 .iter()
519 .any(|endpoint| Rc::ptr_eq(endpoint, &matching[0].endpoint))
520 {
521 return Err(RuntimeFailure::InvalidResolvedPlan {
522 detail: format!(
523 "request binding `{}:{}:{}` does not reference its provider generation endpoint",
524 planned.consumer_instance(),
525 planned.capability_id(),
526 planned.provider_instance()
527 ),
528 });
529 }
530 }
531 if !descriptor.stream_operations().is_empty() {
532 let matching: Vec<_> = stream_bindings
533 .iter()
534 .filter(|prepared| {
535 prepared.requirement_id() == planned.requirement_id()
536 && prepared.consumer_instance == planned.consumer_instance()
537 && prepared.provider_instance == planned.provider_instance()
538 && prepared.endpoint.capability_id() == planned.capability_id()
539 && prepared.endpoint.descriptor_version() == planned.descriptor_version()
540 })
541 .collect();
542 if matching.len() != 1 {
543 return Err(RuntimeFailure::InvalidResolvedPlan {
544 detail: format!(
545 "Execution Adapters prepared {} stream bindings for `{}:{}:{}`; expected 1",
546 matching.len(),
547 planned.consumer_instance(),
548 planned.capability_id(),
549 planned.provider_instance()
550 ),
551 });
552 }
553 if !provider
554 .stream_endpoints()
555 .iter()
556 .any(|endpoint| Rc::ptr_eq(endpoint, &matching[0].endpoint))
557 {
558 return Err(RuntimeFailure::InvalidResolvedPlan {
559 detail: format!(
560 "stream binding `{}:{}:{}` does not reference its provider generation endpoint",
561 planned.consumer_instance(),
562 planned.capability_id(),
563 planned.provider_instance()
564 ),
565 });
566 }
567 }
568 if !descriptor.event_operations().is_empty() {
569 let matching: Vec<_> = event_bindings
570 .iter()
571 .filter(|prepared| {
572 prepared.requirement_id() == planned.requirement_id()
573 && prepared.consumer_instance == planned.consumer_instance()
574 && prepared.provider_instance == planned.provider_instance()
575 && prepared.endpoint.capability_id() == planned.capability_id()
576 && prepared.endpoint.descriptor_version() == planned.descriptor_version()
577 })
578 .collect();
579 if matching.len() != 1 {
580 return Err(RuntimeFailure::InvalidResolvedPlan {
581 detail: format!(
582 "Execution Adapters prepared {} Event bindings for `{}:{}:{}`; expected 1",
583 matching.len(),
584 planned.consumer_instance(),
585 planned.capability_id(),
586 planned.provider_instance()
587 ),
588 });
589 }
590 if !provider
591 .event_endpoints()
592 .iter()
593 .any(|endpoint| Rc::ptr_eq(endpoint, &matching[0].endpoint))
594 {
595 return Err(RuntimeFailure::InvalidResolvedPlan {
596 detail: format!(
597 "Event binding `{}:{}:{}` does not reference its provider generation endpoint",
598 planned.consumer_instance(),
599 planned.capability_id(),
600 planned.provider_instance()
601 ),
602 });
603 }
604 }
605 }
606 Ok(())
607}
608
609pub(super) fn native_plugin_runtimes<D: RuntimeDriver>(
610 plan: &ResolvedAppPlan,
611 driver: &D,
612 mut generations: BTreeMap<String, PreparedNativePlugin>,
613 startup_cancellation: Option<&CancellationToken>,
614) -> BTreeMap<String, NativePluginRuntime> {
615 let mut runtimes = BTreeMap::new();
616 for instance in plan.plugin_instances() {
617 let lifecycle = generations
618 .remove(instance.instance_key())
619 .map(|generation| generation.lifecycle())
620 .expect("prepared App validation requires one generation per planned Instance");
621 let cancellation =
622 startup_cancellation.map_or_else(CancellationToken::new, CancellationToken::child);
623 runtimes.insert(
624 instance.instance_key().to_owned(),
625 NativePluginRuntime {
626 generation: RefCell::new(Some(NativePluginGeneration {
627 lifecycle,
628 tasks: ManagedTaskScope::new_with_cancellation(driver, cancellation),
629 resources: ManagedResourceScope::new(),
630 stop_attempted: false,
631 cleanup_timed_out: false,
632 })),
633 },
634 );
635 }
636 runtimes
637}
638
639pub(super) async fn prepare_native_plugins(
640 runtime: &Rc<NativeAppRuntime>,
641) -> Result<Vec<String>, RuntimeFailure> {
642 let mut prepared_instances = Vec::with_capacity(runtime.activation_order.len());
643 for instance_key in &runtime.activation_order {
644 startup_active(runtime)?;
645 let instance = runtime
646 .plan
647 .plugin_instances()
648 .iter()
649 .find(|instance| instance.instance_key() == instance_key)
650 .expect("activation order only contains planned Plugin Instances");
651 let plugin = runtime
652 .plugins
653 .get(instance_key)
654 .expect("activation order only contains planned Plugin Instances");
655 let (lifecycle, tasks, resources) = plugin
656 .generation_parts()
657 .expect("every startup Plugin Instance has a generation");
658 let cancellation = lifecycle_cancellation(runtime, &tasks);
659 prepared_instances.push(instance_key.clone());
660 let started_at = (runtime.driver.now)();
661 runtime
662 .diagnostics
663 .emit(super::DiagnosticSource::Lifecycle, started_at, |_| {
664 super::DiagnosticEvent::LifecycleStarted {
665 instance: instance_key.clone(),
666 generation: 1,
667 phase: super::PluginLifecyclePhase::Prepare,
668 }
669 });
670 let context = PrepareContext {
671 instance_key: instance_key.clone(),
672 entrypoint: instance.entrypoint().to_owned(),
673 configuration: instance.configuration().to_owned(),
674 dependencies: runtime
675 .dependencies
676 .get(instance_key)
677 .cloned()
678 .unwrap_or_default(),
679 resources,
680 cancellation,
681 admission: runtime.admission.clone(),
682 };
683 let result = lifecycle
684 .prepare(context)
685 .await
686 .and_then(|()| startup_active(runtime));
687 let outcome = result.as_ref().map_or_else(
688 |error| super::DiagnosticOutcome::RuntimeFailure(error.into()),
689 |()| super::DiagnosticOutcome::Succeeded,
690 );
691 runtime.diagnostics.emit(
692 super::DiagnosticSource::Lifecycle,
693 (runtime.driver.now)(),
694 |_| super::DiagnosticEvent::LifecycleCompleted {
695 instance: instance_key.clone(),
696 generation: 1,
697 phase: super::PluginLifecyclePhase::Prepare,
698 outcome,
699 elapsed: (runtime.driver.now)().saturating_sub(started_at),
700 },
701 );
702 if let Err(error) = result {
703 let cleanup_error = deactivate_in_reverse(
704 &runtime.plugins,
705 &runtime.dependencies,
706 &prepared_instances,
707 DeactivationReason::StartupRollback,
708 &runtime.admission,
709 &runtime.diagnostics,
710 &runtime.driver,
711 runtime
712 .startup_cleanup
713 .as_ref()
714 .map(super::cleanup::StartupCleanupBudget::establish),
715 )
716 .await;
717 retain_unsafe_startup(runtime, cleanup_error.as_ref());
718 runtime.diagnostics.emit_runtime_failure(
719 (runtime.driver.now)(),
720 Some(instance_key),
721 &error,
722 );
723 return Err(error);
724 }
725 }
726 Ok(prepared_instances)
727}
728
729fn retain_unsafe_startup(runtime: &Rc<NativeAppRuntime>, cleanup_error: Option<&RuntimeFailure>) {
730 if matches!(cleanup_error, Some(RuntimeFailure::DeadlineExceeded { .. })) {
731 std::mem::forget(runtime.clone());
735 }
736}
737
738pub(super) async fn activate_native_plugins(
739 runtime: &Rc<NativeAppRuntime>,
740) -> Result<(), RuntimeFailure> {
741 for instance_key in &runtime.activation_order {
742 startup_active(runtime)?;
743 let plugin = runtime
744 .plugins
745 .get(instance_key)
746 .expect("activation order only contains planned Plugin Instances");
747 let (lifecycle, tasks, resources) = plugin
748 .generation_parts()
749 .expect("every startup Plugin Instance has a generation");
750 let cancellation = lifecycle_cancellation(runtime, &tasks);
751 let started_at = (runtime.driver.now)();
752 runtime
753 .diagnostics
754 .emit(super::DiagnosticSource::Lifecycle, started_at, |_| {
755 super::DiagnosticEvent::LifecycleStarted {
756 instance: instance_key.clone(),
757 generation: 1,
758 phase: super::PluginLifecyclePhase::Activate,
759 }
760 });
761 let context = ActivateContext {
762 instance_key: instance_key.clone(),
763 dependencies: runtime
764 .dependencies
765 .get(instance_key)
766 .cloned()
767 .unwrap_or_default(),
768 ready_gate: runtime.ready_gate.clone(),
769 tasks,
770 resources,
771 cancellation,
772 admission: runtime.admission.clone(),
773 };
774 let result = lifecycle
775 .activate(context)
776 .await
777 .and_then(|()| startup_active(runtime));
778 let outcome = result.as_ref().map_or_else(
779 |error| super::DiagnosticOutcome::RuntimeFailure(error.into()),
780 |()| super::DiagnosticOutcome::Succeeded,
781 );
782 runtime.diagnostics.emit(
783 super::DiagnosticSource::Lifecycle,
784 (runtime.driver.now)(),
785 |_| super::DiagnosticEvent::LifecycleCompleted {
786 instance: instance_key.clone(),
787 generation: 1,
788 phase: super::PluginLifecyclePhase::Activate,
789 outcome,
790 elapsed: (runtime.driver.now)().saturating_sub(started_at),
791 },
792 );
793 if let Err(error) = result {
794 runtime.diagnostics.emit_runtime_failure(
795 (runtime.driver.now)(),
796 Some(instance_key),
797 &error,
798 );
799 return Err(error);
800 }
801 }
802 Ok(())
803}
804
805pub(super) async fn construct_native_plugins(
806 runtime: &Rc<NativeAppRuntime>,
807) -> Result<(), RuntimeFailure> {
808 for instance_key in &runtime.activation_order {
809 let instance = runtime
810 .plan
811 .plugin_instance(instance_key)
812 .expect("construction order contains planned Instances");
813 if instance.authoring_version() == 1 {
814 continue;
815 }
816 startup_active(runtime)?;
817 let plugin = runtime
818 .plugins
819 .get(instance_key)
820 .expect("construction order contains planned Instances");
821 let (lifecycle, tasks, resources) = plugin
822 .generation_parts()
823 .expect("startup generation exists");
824 let started_at = (runtime.driver.now)();
825 runtime
826 .diagnostics
827 .emit(super::DiagnosticSource::Lifecycle, started_at, |_| {
828 super::DiagnosticEvent::LifecycleStarted {
829 instance: instance_key.clone(),
830 generation: 1,
831 phase: super::PluginLifecyclePhase::Construct,
832 }
833 });
834 let result = lifecycle
835 .construct(ActivateContext {
836 instance_key: instance_key.clone(),
837 dependencies: runtime
838 .dependencies
839 .get(instance_key)
840 .cloned()
841 .unwrap_or_default(),
842 ready_gate: runtime.ready_gate.clone(),
843 tasks: tasks.clone(),
844 resources,
845 cancellation: lifecycle_cancellation(runtime, &tasks),
846 admission: runtime.admission.clone(),
847 })
848 .await
849 .and_then(|()| startup_active(runtime));
850 let outcome = result.as_ref().map_or_else(
851 |error| super::DiagnosticOutcome::RuntimeFailure(error.into()),
852 |()| super::DiagnosticOutcome::Succeeded,
853 );
854 runtime.diagnostics.emit(
855 super::DiagnosticSource::Lifecycle,
856 (runtime.driver.now)(),
857 |_| super::DiagnosticEvent::LifecycleCompleted {
858 instance: instance_key.clone(),
859 generation: 1,
860 phase: super::PluginLifecyclePhase::Construct,
861 outcome,
862 elapsed: (runtime.driver.now)().saturating_sub(started_at),
863 },
864 );
865 result?;
866 }
867 Ok(())
868}
869
870pub(super) async fn open_native_readiness(runtime: &Rc<NativeAppRuntime>) {
871 runtime.startup_context.borrow_mut().take();
872 runtime.ready_gate.open();
873 runtime.admission.open();
874 runtime.diagnostics.emit(
875 super::DiagnosticSource::Lifecycle,
876 (runtime.driver.now)(),
877 |_| super::DiagnosticEvent::AppReady,
878 );
879 (runtime.driver.yield_now)().await;
880}
881
882pub(super) fn native_bindings(
883 plan: &ResolvedAppPlan,
884 prepared: &[PreparedBinding],
885) -> (NativeBindingTable, NativeEndpointStateTable) {
886 let mut bindings = BTreeMap::new();
887 let mut endpoint_states = BTreeMap::new();
888 for binding in plan.capability_bindings() {
889 let Some(descriptor) =
890 plan.plugin_instance(binding.provider_instance())
891 .and_then(|provider| {
892 provider
893 .provided_capabilities()
894 .iter()
895 .find(|endpoint| endpoint.capability_id() == binding.capability_id())
896 })
897 else {
898 continue;
899 };
900 if descriptor.request_operations().is_empty() {
901 continue;
902 }
903 let Some(endpoint) = prepared.iter().find_map(|prepared| {
904 (prepared.requirement_id() == binding.requirement_id()
905 && prepared.consumer_instance == binding.consumer_instance()
906 && prepared.provider_instance == binding.provider_instance()
907 && prepared.endpoint.capability_id() == binding.capability_id())
908 .then_some(&prepared.endpoint)
909 }) else {
910 continue;
911 };
912 let state = endpoint_states
913 .entry((
914 binding.provider_instance().to_owned(),
915 endpoint.capability_id().to_owned(),
916 ))
917 .or_insert_with(|| Rc::new(NativeEndpointState::new(endpoint.clone(), 1)))
918 .clone();
919 let admissions = endpoint
920 .operations()
921 .iter()
922 .map(|operation| {
923 (
924 (*operation).to_owned(),
925 RequestAdmission::new(plan.request_admission_for(binding, operation)),
926 )
927 })
928 .collect();
929 bindings
930 .entry((
931 binding.consumer_instance().to_owned(),
932 endpoint.capability_id(),
933 ))
934 .or_insert_with(Vec::new)
935 .push(NativeEndpointBinding {
936 requirement_id: binding.requirement_id().to_owned(),
937 plugin_instance: binding.provider_instance().to_owned(),
938 state,
939 admissions,
940 });
941 }
942 (bindings, endpoint_states)
943}
944
945pub(super) fn native_stream_bindings(
946 plan: &ResolvedAppPlan,
947 prepared: &[PreparedStreamBinding],
948) -> (NativeStreamBindingTable, NativeStreamEndpointStateTable) {
949 let mut bindings = BTreeMap::new();
950 let mut endpoint_states = BTreeMap::new();
951 for binding in plan.capability_bindings() {
952 let Some(descriptor) =
953 plan.plugin_instance(binding.provider_instance())
954 .and_then(|provider| {
955 provider
956 .provided_capabilities()
957 .iter()
958 .find(|endpoint| endpoint.capability_id() == binding.capability_id())
959 })
960 else {
961 continue;
962 };
963 if descriptor.stream_operations().is_empty() {
964 continue;
965 }
966 let Some(endpoint) = prepared.iter().find_map(|prepared| {
967 (prepared.requirement_id() == binding.requirement_id()
968 && prepared.consumer_instance == binding.consumer_instance()
969 && prepared.provider_instance == binding.provider_instance()
970 && prepared.endpoint.capability_id() == binding.capability_id())
971 .then_some(&prepared.endpoint)
972 }) else {
973 continue;
974 };
975 let state = endpoint_states
976 .entry((
977 binding.provider_instance().to_owned(),
978 endpoint.capability_id().to_owned(),
979 ))
980 .or_insert_with(|| Rc::new(NativeStreamEndpointState::new(endpoint.clone(), 1)))
981 .clone();
982 let admissions = endpoint
983 .operations()
984 .iter()
985 .map(|operation| {
986 (
987 (*operation).to_owned(),
988 RequestAdmission::new(plan.request_admission_for(binding, operation)),
989 )
990 })
991 .collect();
992 bindings
993 .entry((
994 binding.consumer_instance().to_owned(),
995 endpoint.capability_id(),
996 ))
997 .or_insert_with(Vec::new)
998 .push(NativeStreamEndpointBinding {
999 requirement_id: binding.requirement_id().to_owned(),
1000 plugin_instance: binding.provider_instance().to_owned(),
1001 state,
1002 admissions,
1003 });
1004 }
1005 (bindings, endpoint_states)
1006}
1007
1008pub(super) fn native_event_bindings(
1009 plan: &ResolvedAppPlan,
1010 prepared: &[PreparedEventBinding],
1011) -> (NativeEventBindingTable, NativeEventEndpointStateTable) {
1012 let mut bindings = BTreeMap::new();
1013 let mut endpoint_states = BTreeMap::new();
1014 for binding in plan.capability_bindings() {
1015 let Some(descriptor) =
1016 plan.plugin_instance(binding.provider_instance())
1017 .and_then(|provider| {
1018 provider
1019 .provided_capabilities()
1020 .iter()
1021 .find(|endpoint| endpoint.capability_id() == binding.capability_id())
1022 })
1023 else {
1024 continue;
1025 };
1026 if descriptor.event_operations().is_empty() {
1027 continue;
1028 }
1029 let Some(endpoint) = prepared.iter().find_map(|prepared| {
1030 (prepared.requirement_id() == binding.requirement_id()
1031 && prepared.consumer_instance == binding.consumer_instance()
1032 && prepared.provider_instance == binding.provider_instance()
1033 && prepared.endpoint.capability_id() == binding.capability_id())
1034 .then_some(&prepared.endpoint)
1035 }) else {
1036 continue;
1037 };
1038 let state = endpoint_states
1039 .entry((
1040 binding.provider_instance().to_owned(),
1041 endpoint.capability_id().to_owned(),
1042 ))
1043 .or_insert_with(|| Rc::new(event::NativeEventEndpointState::new(endpoint.clone(), 1)))
1044 .clone();
1045 let queue = event::NativeEventQueue::new(plan.event_admission_for(binding));
1046 state.register_queue(&queue);
1047 bindings
1048 .entry((
1049 binding.consumer_instance().to_owned(),
1050 endpoint.capability_id(),
1051 ))
1052 .or_insert_with(Vec::new)
1053 .push(event::NativeEventEndpointBinding {
1054 requirement_id: binding.requirement_id().to_owned(),
1055 plugin_instance: binding.provider_instance().to_owned(),
1056 state,
1057 queue,
1058 });
1059 }
1060 (bindings, endpoint_states)
1061}
1062
1063pub(super) fn plugin_dependencies(
1064 plan: &ResolvedAppPlan,
1065 endpoints: &BTreeMap<(String, &'static str), Vec<NativeEndpointBinding>>,
1066 stream_endpoints: &NativeStreamBindingTable,
1067 event_endpoints: &NativeEventBindingTable,
1068 runtime: &Rc<RefCell<Weak<NativeAppRuntime>>>,
1069) -> BTreeMap<String, PluginDependencies> {
1070 let mut dependencies: BTreeMap<String, PluginDependencies> = plan
1071 .plugin_instances()
1072 .iter()
1073 .map(|instance| {
1074 (
1075 instance.instance_key().to_owned(),
1076 PluginDependencies::new(
1077 instance.instance_key(),
1078 runtime.clone(),
1079 instance.required_capabilities().to_vec(),
1080 ),
1081 )
1082 })
1083 .collect();
1084 for binding in plan.capability_bindings() {
1085 dependencies
1086 .get_mut(binding.consumer_instance())
1087 .expect("every resolved binding consumer has Plugin dependencies")
1088 .bindings
1089 .push(PluginDependency::new(
1090 binding.requirement_id(),
1091 binding.capability_id(),
1092 binding.provider_instance(),
1093 binding.provider_order(),
1094 endpoints
1095 .iter()
1096 .find(|((consumer, capability), _)| {
1097 consumer == binding.consumer_instance()
1098 && *capability == binding.capability_id()
1099 })
1100 .and_then(|(_, endpoints)| {
1101 endpoints.iter().find(|endpoint| {
1102 endpoint.requirement_id == binding.requirement_id()
1103 && endpoint.plugin_instance == binding.provider_instance()
1104 })
1105 })
1106 .map(|endpoint| PluginDependencyHandle {
1107 binding: endpoint.clone(),
1108 caller_instance: binding.consumer_instance().to_owned(),
1109 runtime: runtime.clone(),
1110 }),
1111 stream_endpoints
1112 .iter()
1113 .find(|((consumer, capability), _)| {
1114 consumer == binding.consumer_instance()
1115 && *capability == binding.capability_id()
1116 })
1117 .and_then(|(_, endpoints)| {
1118 endpoints.iter().find(|endpoint| {
1119 endpoint.requirement_id == binding.requirement_id()
1120 && endpoint.plugin_instance == binding.provider_instance()
1121 })
1122 })
1123 .map(|endpoint| PluginStreamDependencyHandle {
1124 binding: endpoint.clone(),
1125 caller_instance: binding.consumer_instance().to_owned(),
1126 runtime: runtime.clone(),
1127 }),
1128 event_endpoints
1129 .iter()
1130 .find(|((consumer, capability), _)| {
1131 consumer == binding.consumer_instance()
1132 && *capability == binding.capability_id()
1133 })
1134 .and_then(|(_, endpoints)| {
1135 endpoints.iter().find(|endpoint| {
1136 endpoint.requirement_id == binding.requirement_id()
1137 && endpoint.plugin_instance == binding.provider_instance()
1138 })
1139 })
1140 .map(|endpoint| PluginEventDependencyHandle {
1141 binding: endpoint.clone(),
1142 caller_instance: binding.consumer_instance().to_owned(),
1143 runtime: runtime.clone(),
1144 }),
1145 ));
1146 }
1147 dependencies
1148}