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 cleanup_failure: RefCell::new(None),
227 });
228 runtime_link.replace(Rc::downgrade(&runtime));
229 attach_managed_task_failure_handlers(&runtime);
230 runtime.diagnostics.emit(
231 super::DiagnosticSource::Lifecycle,
232 (runtime.driver.now)(),
233 |_| super::DiagnosticEvent::AppStarted {
234 plugin_count: runtime.plan.plugin_instances().len(),
235 },
236 );
237 let prepared_instances = prepare_native_plugins(&runtime).await?;
238 if let Err(error) = activate_native_plugins(&runtime).await {
239 let cleanup_error = deactivate_in_reverse(
240 &runtime.plugins,
241 &runtime.dependencies,
242 &prepared_instances,
243 DeactivationReason::StartupRollback,
244 &runtime.admission,
245 &runtime.diagnostics,
246 &runtime.driver,
247 runtime
248 .startup_cleanup
249 .as_ref()
250 .map(super::cleanup::StartupCleanupBudget::establish),
251 )
252 .await;
253 retain_unsafe_startup(&runtime, cleanup_error.as_ref());
254 runtime
255 .diagnostics
256 .emit_runtime_failure((runtime.driver.now)(), None, &error);
257 return Err(error);
258 }
259 open_native_readiness(&runtime).await;
260 Ok(NativeApp {
261 bindings,
262 stream_bindings,
263 event_bindings,
264 diagnostics,
265 runtime,
266 })
267 }
268}
269
270pub(super) fn attach_managed_task_failure_handlers(runtime: &Rc<NativeAppRuntime>) {
271 for (instance_key, plugin) in &runtime.plugins {
272 let Some((_, tasks, _)) = plugin.generation_parts() else {
273 continue;
274 };
275 attach_managed_task_failure_handler(runtime, instance_key, &tasks);
276 }
277}
278
279fn startup_active(runtime: &NativeAppRuntime) -> Result<(), RuntimeFailure> {
280 let result = runtime
281 .startup_context
282 .borrow()
283 .as_ref()
284 .map_or(Ok(()), |context| {
285 super::ensure_context_active(&runtime.driver, context)
286 });
287 if result.is_err()
288 && let Some(cleanup) = &runtime.startup_cleanup
289 {
290 let now = (runtime.driver.now)();
291 let cleanup_started_at = runtime
292 .startup_context
293 .borrow()
294 .as_ref()
295 .and_then(super::InvocationContext::deadline)
296 .filter(|deadline| now >= *deadline)
297 .unwrap_or(now);
298 cleanup.establish_at(cleanup_started_at);
299 }
300 result
301}
302
303fn lifecycle_cancellation(
304 _runtime: &NativeAppRuntime,
305 tasks: &ManagedTaskScope,
306) -> CancellationToken {
307 tasks.cancellation()
308}
309
310pub(super) fn attach_managed_task_failure_handler(
311 runtime: &Rc<NativeAppRuntime>,
312 instance_key: &str,
313 tasks: &ManagedTaskScope,
314) {
315 let task_runtime = Rc::downgrade(runtime);
316 let task_instance_key = instance_key.to_owned();
317 let handler: Rc<dyn Fn()> = Rc::new(move || {
318 let Some(runtime) = task_runtime.upgrade() else {
319 return;
320 };
321 if begin_plugin_supervision(&runtime, &task_instance_key).unwrap_or(false)
322 && let Err(error) = schedule_plugin_supervision(&runtime, &task_instance_key)
323 {
324 let _ = handle_supervision_schedule_failure(&runtime, &task_instance_key, error);
325 }
326 });
327 tasks.set_failure_handler(&handler);
328}
329
330pub(super) fn runtime_plan_error(error: &PlanResolutionError) -> RuntimeFailure {
331 RuntimeFailure::InvalidResolvedPlan {
332 detail: error.to_string(),
333 }
334}
335
336#[allow(
337 clippy::too_many_lines,
338 reason = "one fail-closed pass keeps request, stream, event, and generation validation aligned"
339)]
340pub(super) fn validate_prepared_native_app(
341 plan: &ResolvedAppPlan,
342 bindings: &[PreparedBinding],
343 stream_bindings: &[PreparedStreamBinding],
344 event_bindings: &[PreparedEventBinding],
345 generations: &BTreeMap<String, PreparedNativePlugin>,
346) -> Result<(), RuntimeFailure> {
347 if generations.len() != plan.plugin_instances().len() {
348 return Err(RuntimeFailure::InvalidResolvedPlan {
349 detail: format!(
350 "Execution Adapters prepared {} Plugin generations; expected {}",
351 generations.len(),
352 plan.plugin_instances().len()
353 ),
354 });
355 }
356 for instance in plan.plugin_instances() {
357 let generation = generations.get(instance.instance_key()).ok_or_else(|| {
358 RuntimeFailure::InvalidResolvedPlan {
359 detail: format!(
360 "Execution Adapters did not prepare Plugin Instance `{}`",
361 instance.instance_key()
362 ),
363 }
364 })?;
365 validate_native_endpoint_set(
366 instance.instance_key(),
367 instance,
368 generation.endpoints(),
369 generation.stream_endpoints(),
370 generation.event_endpoints(),
371 )?;
372 }
373 if let Some(instance_key) = generations.keys().find(|instance_key| {
374 !plan
375 .plugin_instances()
376 .iter()
377 .any(|instance| instance.instance_key() == instance_key.as_str())
378 }) {
379 return Err(RuntimeFailure::InvalidResolvedPlan {
380 detail: format!("Execution Adapter prepared unknown Plugin Instance `{instance_key}`"),
381 });
382 }
383
384 let expected_request_bindings = plan
385 .capability_bindings()
386 .iter()
387 .filter(|binding| {
388 plan.plugin_instance(binding.provider_instance())
389 .and_then(|provider| {
390 provider
391 .provided_capabilities()
392 .iter()
393 .find(|endpoint| endpoint.capability_id() == binding.capability_id())
394 })
395 .is_some_and(|endpoint| !endpoint.request_operations().is_empty())
396 })
397 .count();
398 let expected_stream_bindings = plan
399 .capability_bindings()
400 .iter()
401 .filter(|binding| {
402 plan.plugin_instance(binding.provider_instance())
403 .and_then(|provider| {
404 provider
405 .provided_capabilities()
406 .iter()
407 .find(|endpoint| endpoint.capability_id() == binding.capability_id())
408 })
409 .is_some_and(|endpoint| !endpoint.stream_operations().is_empty())
410 })
411 .count();
412 let expected_event_bindings = plan
413 .capability_bindings()
414 .iter()
415 .filter(|binding| {
416 plan.plugin_instance(binding.provider_instance())
417 .and_then(|provider| {
418 provider
419 .provided_capabilities()
420 .iter()
421 .find(|endpoint| endpoint.capability_id() == binding.capability_id())
422 })
423 .is_some_and(|endpoint| !endpoint.event_operations().is_empty())
424 })
425 .count();
426 if bindings.len() != expected_request_bindings {
427 return Err(RuntimeFailure::InvalidResolvedPlan {
428 detail: if expected_stream_bindings == 0 && stream_bindings.is_empty() {
429 format!(
430 "Execution Adapters prepared {} bindings; expected {}",
431 bindings.len(),
432 expected_request_bindings
433 )
434 } else {
435 format!(
436 "Execution Adapters prepared {} request bindings; expected {}",
437 bindings.len(),
438 expected_request_bindings
439 )
440 },
441 });
442 }
443 if stream_bindings.len() != expected_stream_bindings {
444 return Err(RuntimeFailure::InvalidResolvedPlan {
445 detail: format!(
446 "Execution Adapters prepared {} stream bindings; expected {}",
447 stream_bindings.len(),
448 expected_stream_bindings
449 ),
450 });
451 }
452 if event_bindings.len() != expected_event_bindings {
453 return Err(RuntimeFailure::InvalidResolvedPlan {
454 detail: format!(
455 "Execution Adapters prepared {} Event bindings; expected {}",
456 event_bindings.len(),
457 expected_event_bindings
458 ),
459 });
460 }
461 for planned in plan.capability_bindings() {
462 let provider = generations
463 .get(planned.provider_instance())
464 .expect("the resolved Plan references one validated provider generation");
465 let descriptor = plan
466 .plugin_instance(planned.provider_instance())
467 .and_then(|provider| {
468 provider
469 .provided_capabilities()
470 .iter()
471 .find(|endpoint| endpoint.capability_id() == planned.capability_id())
472 })
473 .expect("the resolved Plan references one validated provider endpoint");
474 if !descriptor.request_operations().is_empty() {
475 let matching: Vec<_> = bindings
476 .iter()
477 .filter(|prepared| {
478 prepared.requirement_id() == planned.requirement_id()
479 && prepared.consumer_instance == planned.consumer_instance()
480 && prepared.provider_instance == planned.provider_instance()
481 && prepared.endpoint.capability_id() == planned.capability_id()
482 && prepared.endpoint.descriptor_version() == planned.descriptor_version()
483 })
484 .collect();
485 if matching.len() != 1 {
486 return Err(RuntimeFailure::InvalidResolvedPlan {
487 detail: format!(
488 "Execution Adapters prepared {} request bindings for `{}:{}:{}`; expected 1",
489 matching.len(),
490 planned.consumer_instance(),
491 planned.capability_id(),
492 planned.provider_instance()
493 ),
494 });
495 }
496 if !provider
497 .endpoints()
498 .iter()
499 .any(|endpoint| Rc::ptr_eq(endpoint, &matching[0].endpoint))
500 {
501 return Err(RuntimeFailure::InvalidResolvedPlan {
502 detail: format!(
503 "request binding `{}:{}:{}` does not reference its provider generation endpoint",
504 planned.consumer_instance(),
505 planned.capability_id(),
506 planned.provider_instance()
507 ),
508 });
509 }
510 }
511 if !descriptor.stream_operations().is_empty() {
512 let matching: Vec<_> = stream_bindings
513 .iter()
514 .filter(|prepared| {
515 prepared.requirement_id() == planned.requirement_id()
516 && prepared.consumer_instance == planned.consumer_instance()
517 && prepared.provider_instance == planned.provider_instance()
518 && prepared.endpoint.capability_id() == planned.capability_id()
519 && prepared.endpoint.descriptor_version() == planned.descriptor_version()
520 })
521 .collect();
522 if matching.len() != 1 {
523 return Err(RuntimeFailure::InvalidResolvedPlan {
524 detail: format!(
525 "Execution Adapters prepared {} stream bindings for `{}:{}:{}`; expected 1",
526 matching.len(),
527 planned.consumer_instance(),
528 planned.capability_id(),
529 planned.provider_instance()
530 ),
531 });
532 }
533 if !provider
534 .stream_endpoints()
535 .iter()
536 .any(|endpoint| Rc::ptr_eq(endpoint, &matching[0].endpoint))
537 {
538 return Err(RuntimeFailure::InvalidResolvedPlan {
539 detail: format!(
540 "stream binding `{}:{}:{}` does not reference its provider generation endpoint",
541 planned.consumer_instance(),
542 planned.capability_id(),
543 planned.provider_instance()
544 ),
545 });
546 }
547 }
548 if !descriptor.event_operations().is_empty() {
549 let matching: Vec<_> = event_bindings
550 .iter()
551 .filter(|prepared| {
552 prepared.requirement_id() == planned.requirement_id()
553 && prepared.consumer_instance == planned.consumer_instance()
554 && prepared.provider_instance == planned.provider_instance()
555 && prepared.endpoint.capability_id() == planned.capability_id()
556 && prepared.endpoint.descriptor_version() == planned.descriptor_version()
557 })
558 .collect();
559 if matching.len() != 1 {
560 return Err(RuntimeFailure::InvalidResolvedPlan {
561 detail: format!(
562 "Execution Adapters prepared {} Event bindings for `{}:{}:{}`; expected 1",
563 matching.len(),
564 planned.consumer_instance(),
565 planned.capability_id(),
566 planned.provider_instance()
567 ),
568 });
569 }
570 if !provider
571 .event_endpoints()
572 .iter()
573 .any(|endpoint| Rc::ptr_eq(endpoint, &matching[0].endpoint))
574 {
575 return Err(RuntimeFailure::InvalidResolvedPlan {
576 detail: format!(
577 "Event binding `{}:{}:{}` does not reference its provider generation endpoint",
578 planned.consumer_instance(),
579 planned.capability_id(),
580 planned.provider_instance()
581 ),
582 });
583 }
584 }
585 }
586 Ok(())
587}
588
589pub(super) fn native_plugin_runtimes<D: RuntimeDriver>(
590 plan: &ResolvedAppPlan,
591 driver: &D,
592 mut generations: BTreeMap<String, PreparedNativePlugin>,
593 startup_cancellation: Option<&CancellationToken>,
594) -> BTreeMap<String, NativePluginRuntime> {
595 let mut runtimes = BTreeMap::new();
596 for instance in plan.plugin_instances() {
597 let lifecycle = generations
598 .remove(instance.instance_key())
599 .map(|generation| generation.lifecycle())
600 .expect("prepared App validation requires one generation per planned Instance");
601 let cancellation =
602 startup_cancellation.map_or_else(CancellationToken::new, CancellationToken::child);
603 runtimes.insert(
604 instance.instance_key().to_owned(),
605 NativePluginRuntime {
606 generation: RefCell::new(Some(NativePluginGeneration {
607 lifecycle,
608 tasks: ManagedTaskScope::new_with_cancellation(driver, cancellation),
609 resources: ManagedResourceScope::new(),
610 stop_attempted: false,
611 cleanup_timed_out: false,
612 })),
613 },
614 );
615 }
616 runtimes
617}
618
619pub(super) async fn prepare_native_plugins(
620 runtime: &Rc<NativeAppRuntime>,
621) -> Result<Vec<String>, RuntimeFailure> {
622 let mut prepared_instances = Vec::with_capacity(runtime.activation_order.len());
623 for instance_key in &runtime.activation_order {
624 startup_active(runtime)?;
625 let instance = runtime
626 .plan
627 .plugin_instances()
628 .iter()
629 .find(|instance| instance.instance_key() == instance_key)
630 .expect("activation order only contains planned Plugin Instances");
631 let plugin = runtime
632 .plugins
633 .get(instance_key)
634 .expect("activation order only contains planned Plugin Instances");
635 let (lifecycle, tasks, resources) = plugin
636 .generation_parts()
637 .expect("every startup Plugin Instance has a generation");
638 let cancellation = lifecycle_cancellation(runtime, &tasks);
639 prepared_instances.push(instance_key.clone());
640 let started_at = (runtime.driver.now)();
641 runtime
642 .diagnostics
643 .emit(super::DiagnosticSource::Lifecycle, started_at, |_| {
644 super::DiagnosticEvent::LifecycleStarted {
645 instance: instance_key.clone(),
646 generation: 1,
647 phase: super::PluginLifecyclePhase::Prepare,
648 }
649 });
650 let context = PrepareContext {
651 instance_key: instance_key.clone(),
652 entrypoint: instance.entrypoint().to_owned(),
653 configuration: instance.configuration().to_owned(),
654 dependencies: runtime
655 .dependencies
656 .get(instance_key)
657 .cloned()
658 .unwrap_or_default(),
659 resources,
660 cancellation,
661 admission: runtime.admission.clone(),
662 };
663 let result = lifecycle
664 .prepare(context)
665 .await
666 .and_then(|()| startup_active(runtime));
667 let outcome = result.as_ref().map_or_else(
668 |error| super::DiagnosticOutcome::RuntimeFailure(error.into()),
669 |()| super::DiagnosticOutcome::Succeeded,
670 );
671 runtime.diagnostics.emit(
672 super::DiagnosticSource::Lifecycle,
673 (runtime.driver.now)(),
674 |_| super::DiagnosticEvent::LifecycleCompleted {
675 instance: instance_key.clone(),
676 generation: 1,
677 phase: super::PluginLifecyclePhase::Prepare,
678 outcome,
679 elapsed: (runtime.driver.now)().saturating_sub(started_at),
680 },
681 );
682 let result = match result {
683 Ok(()) => construct_native_plugin(runtime, instance_key).await,
684 Err(error) => Err(error),
685 };
686 if let Err(error) = result {
687 let cleanup_error = deactivate_in_reverse(
688 &runtime.plugins,
689 &runtime.dependencies,
690 &prepared_instances,
691 DeactivationReason::StartupRollback,
692 &runtime.admission,
693 &runtime.diagnostics,
694 &runtime.driver,
695 runtime
696 .startup_cleanup
697 .as_ref()
698 .map(super::cleanup::StartupCleanupBudget::establish),
699 )
700 .await;
701 retain_unsafe_startup(runtime, cleanup_error.as_ref());
702 runtime.diagnostics.emit_runtime_failure(
703 (runtime.driver.now)(),
704 Some(instance_key),
705 &error,
706 );
707 return Err(error);
708 }
709 }
710 Ok(prepared_instances)
711}
712
713fn retain_unsafe_startup(runtime: &Rc<NativeAppRuntime>, cleanup_error: Option<&RuntimeFailure>) {
714 if let Some(error) = cleanup_error {
715 runtime.record_cleanup_failure(error);
716 }
717 if matches!(cleanup_error, Some(RuntimeFailure::DeadlineExceeded { .. })) {
718 std::mem::forget(runtime.clone());
722 }
723}
724
725pub(super) async fn activate_native_plugins(
726 runtime: &Rc<NativeAppRuntime>,
727) -> Result<(), RuntimeFailure> {
728 for instance_key in &runtime.activation_order {
729 startup_active(runtime)?;
730 let plugin = runtime
731 .plugins
732 .get(instance_key)
733 .expect("activation order only contains planned Plugin Instances");
734 let (lifecycle, tasks, resources) = plugin
735 .generation_parts()
736 .expect("every startup Plugin Instance has a generation");
737 let cancellation = lifecycle_cancellation(runtime, &tasks);
738 let started_at = (runtime.driver.now)();
739 runtime
740 .diagnostics
741 .emit(super::DiagnosticSource::Lifecycle, started_at, |_| {
742 super::DiagnosticEvent::LifecycleStarted {
743 instance: instance_key.clone(),
744 generation: 1,
745 phase: super::PluginLifecyclePhase::Activate,
746 }
747 });
748 let context = ActivateContext {
749 instance_key: instance_key.clone(),
750 dependencies: runtime
751 .dependencies
752 .get(instance_key)
753 .cloned()
754 .unwrap_or_default(),
755 ready_gate: runtime.ready_gate.clone(),
756 tasks,
757 resources,
758 cancellation,
759 admission: runtime.admission.clone(),
760 };
761 let result = lifecycle
762 .activate(context)
763 .await
764 .and_then(|()| startup_active(runtime));
765 let outcome = result.as_ref().map_or_else(
766 |error| super::DiagnosticOutcome::RuntimeFailure(error.into()),
767 |()| super::DiagnosticOutcome::Succeeded,
768 );
769 runtime.diagnostics.emit(
770 super::DiagnosticSource::Lifecycle,
771 (runtime.driver.now)(),
772 |_| super::DiagnosticEvent::LifecycleCompleted {
773 instance: instance_key.clone(),
774 generation: 1,
775 phase: super::PluginLifecyclePhase::Activate,
776 outcome,
777 elapsed: (runtime.driver.now)().saturating_sub(started_at),
778 },
779 );
780 if let Err(error) = result {
781 runtime.diagnostics.emit_runtime_failure(
782 (runtime.driver.now)(),
783 Some(instance_key),
784 &error,
785 );
786 return Err(error);
787 }
788 }
789 Ok(())
790}
791
792async fn construct_native_plugin(
793 runtime: &Rc<NativeAppRuntime>,
794 instance_key: &str,
795) -> Result<(), RuntimeFailure> {
796 let instance = runtime
797 .plan
798 .plugin_instance(instance_key)
799 .expect("construction order contains planned Instances");
800 if instance.authoring_version() == 1 {
801 return Ok(());
802 }
803 startup_active(runtime)?;
804 let plugin = runtime
805 .plugins
806 .get(instance_key)
807 .expect("construction order contains planned Instances");
808 let (lifecycle, tasks, resources) = plugin
809 .generation_parts()
810 .expect("startup generation exists");
811 let started_at = (runtime.driver.now)();
812 runtime
813 .diagnostics
814 .emit(super::DiagnosticSource::Lifecycle, started_at, |_| {
815 super::DiagnosticEvent::LifecycleStarted {
816 instance: instance_key.to_owned(),
817 generation: 1,
818 phase: super::PluginLifecyclePhase::Construct,
819 }
820 });
821 let result = lifecycle
822 .construct(ActivateContext {
823 instance_key: instance_key.to_owned(),
824 dependencies: runtime
825 .dependencies
826 .get(instance_key)
827 .cloned()
828 .unwrap_or_default(),
829 ready_gate: runtime.ready_gate.clone(),
830 tasks: tasks.clone(),
831 resources,
832 cancellation: lifecycle_cancellation(runtime, &tasks),
833 admission: runtime.admission.clone(),
834 })
835 .await
836 .and_then(|()| startup_active(runtime));
837 let outcome = result.as_ref().map_or_else(
838 |error| super::DiagnosticOutcome::RuntimeFailure(error.into()),
839 |()| super::DiagnosticOutcome::Succeeded,
840 );
841 runtime.diagnostics.emit(
842 super::DiagnosticSource::Lifecycle,
843 (runtime.driver.now)(),
844 |_| super::DiagnosticEvent::LifecycleCompleted {
845 instance: instance_key.to_owned(),
846 generation: 1,
847 phase: super::PluginLifecyclePhase::Construct,
848 outcome,
849 elapsed: (runtime.driver.now)().saturating_sub(started_at),
850 },
851 );
852 result
853}
854
855pub(super) async fn open_native_readiness(runtime: &Rc<NativeAppRuntime>) {
856 runtime.startup_context.borrow_mut().take();
857 runtime.ready_gate.open();
858 runtime.admission.open();
859 runtime.diagnostics.emit(
860 super::DiagnosticSource::Lifecycle,
861 (runtime.driver.now)(),
862 |_| super::DiagnosticEvent::AppReady,
863 );
864 (runtime.driver.yield_now)().await;
865}
866
867pub(super) fn native_bindings(
868 plan: &ResolvedAppPlan,
869 prepared: &[PreparedBinding],
870) -> (NativeBindingTable, NativeEndpointStateTable) {
871 let mut bindings = BTreeMap::new();
872 let mut endpoint_states = BTreeMap::new();
873 for binding in plan.capability_bindings() {
874 let Some(descriptor) =
875 plan.plugin_instance(binding.provider_instance())
876 .and_then(|provider| {
877 provider
878 .provided_capabilities()
879 .iter()
880 .find(|endpoint| endpoint.capability_id() == binding.capability_id())
881 })
882 else {
883 continue;
884 };
885 if descriptor.request_operations().is_empty() {
886 continue;
887 }
888 let Some(endpoint) = prepared.iter().find_map(|prepared| {
889 (prepared.requirement_id() == binding.requirement_id()
890 && prepared.consumer_instance == binding.consumer_instance()
891 && prepared.provider_instance == binding.provider_instance()
892 && prepared.endpoint.capability_id() == binding.capability_id())
893 .then_some(&prepared.endpoint)
894 }) else {
895 continue;
896 };
897 let state = endpoint_states
898 .entry((
899 binding.provider_instance().to_owned(),
900 endpoint.capability_id().to_owned(),
901 ))
902 .or_insert_with(|| Rc::new(NativeEndpointState::new(endpoint.clone(), 1)))
903 .clone();
904 let admissions = endpoint
905 .operations()
906 .iter()
907 .map(|operation| {
908 (
909 (*operation).to_owned(),
910 RequestAdmission::new(plan.request_admission_for(binding, operation)),
911 )
912 })
913 .collect();
914 bindings
915 .entry((
916 binding.consumer_instance().to_owned(),
917 endpoint.capability_id(),
918 ))
919 .or_insert_with(Vec::new)
920 .push(NativeEndpointBinding {
921 requirement_id: binding.requirement_id().to_owned(),
922 plugin_instance: binding.provider_instance().to_owned(),
923 state,
924 admissions,
925 });
926 }
927 (bindings, endpoint_states)
928}
929
930pub(super) fn native_stream_bindings(
931 plan: &ResolvedAppPlan,
932 prepared: &[PreparedStreamBinding],
933) -> (NativeStreamBindingTable, NativeStreamEndpointStateTable) {
934 let mut bindings = BTreeMap::new();
935 let mut endpoint_states = BTreeMap::new();
936 for binding in plan.capability_bindings() {
937 let Some(descriptor) =
938 plan.plugin_instance(binding.provider_instance())
939 .and_then(|provider| {
940 provider
941 .provided_capabilities()
942 .iter()
943 .find(|endpoint| endpoint.capability_id() == binding.capability_id())
944 })
945 else {
946 continue;
947 };
948 if descriptor.stream_operations().is_empty() {
949 continue;
950 }
951 let Some(endpoint) = prepared.iter().find_map(|prepared| {
952 (prepared.requirement_id() == binding.requirement_id()
953 && prepared.consumer_instance == binding.consumer_instance()
954 && prepared.provider_instance == binding.provider_instance()
955 && prepared.endpoint.capability_id() == binding.capability_id())
956 .then_some(&prepared.endpoint)
957 }) else {
958 continue;
959 };
960 let state = endpoint_states
961 .entry((
962 binding.provider_instance().to_owned(),
963 endpoint.capability_id().to_owned(),
964 ))
965 .or_insert_with(|| Rc::new(NativeStreamEndpointState::new(endpoint.clone(), 1)))
966 .clone();
967 let admissions = endpoint
968 .operations()
969 .iter()
970 .map(|operation| {
971 (
972 (*operation).to_owned(),
973 RequestAdmission::new(plan.request_admission_for(binding, operation)),
974 )
975 })
976 .collect();
977 bindings
978 .entry((
979 binding.consumer_instance().to_owned(),
980 endpoint.capability_id(),
981 ))
982 .or_insert_with(Vec::new)
983 .push(NativeStreamEndpointBinding {
984 requirement_id: binding.requirement_id().to_owned(),
985 plugin_instance: binding.provider_instance().to_owned(),
986 state,
987 admissions,
988 });
989 }
990 (bindings, endpoint_states)
991}
992
993pub(super) fn native_event_bindings(
994 plan: &ResolvedAppPlan,
995 prepared: &[PreparedEventBinding],
996) -> (NativeEventBindingTable, NativeEventEndpointStateTable) {
997 let mut bindings = BTreeMap::new();
998 let mut endpoint_states = BTreeMap::new();
999 for binding in plan.capability_bindings() {
1000 let Some(descriptor) =
1001 plan.plugin_instance(binding.provider_instance())
1002 .and_then(|provider| {
1003 provider
1004 .provided_capabilities()
1005 .iter()
1006 .find(|endpoint| endpoint.capability_id() == binding.capability_id())
1007 })
1008 else {
1009 continue;
1010 };
1011 if descriptor.event_operations().is_empty() {
1012 continue;
1013 }
1014 let Some(endpoint) = prepared.iter().find_map(|prepared| {
1015 (prepared.requirement_id() == binding.requirement_id()
1016 && prepared.consumer_instance == binding.consumer_instance()
1017 && prepared.provider_instance == binding.provider_instance()
1018 && prepared.endpoint.capability_id() == binding.capability_id())
1019 .then_some(&prepared.endpoint)
1020 }) else {
1021 continue;
1022 };
1023 let state = endpoint_states
1024 .entry((
1025 binding.provider_instance().to_owned(),
1026 endpoint.capability_id().to_owned(),
1027 ))
1028 .or_insert_with(|| Rc::new(event::NativeEventEndpointState::new(endpoint.clone(), 1)))
1029 .clone();
1030 let queue = event::NativeEventQueue::new(plan.event_admission_for(binding));
1031 state.register_queue(&queue);
1032 bindings
1033 .entry((
1034 binding.consumer_instance().to_owned(),
1035 endpoint.capability_id(),
1036 ))
1037 .or_insert_with(Vec::new)
1038 .push(event::NativeEventEndpointBinding {
1039 requirement_id: binding.requirement_id().to_owned(),
1040 plugin_instance: binding.provider_instance().to_owned(),
1041 state,
1042 queue,
1043 });
1044 }
1045 (bindings, endpoint_states)
1046}
1047
1048pub(super) fn plugin_dependencies(
1049 plan: &ResolvedAppPlan,
1050 endpoints: &BTreeMap<(String, &'static str), Vec<NativeEndpointBinding>>,
1051 stream_endpoints: &NativeStreamBindingTable,
1052 event_endpoints: &NativeEventBindingTable,
1053 runtime: &Rc<RefCell<Weak<NativeAppRuntime>>>,
1054) -> BTreeMap<String, PluginDependencies> {
1055 let mut dependencies: BTreeMap<String, PluginDependencies> = plan
1056 .plugin_instances()
1057 .iter()
1058 .map(|instance| {
1059 (
1060 instance.instance_key().to_owned(),
1061 PluginDependencies::new(
1062 instance.instance_key(),
1063 runtime.clone(),
1064 instance.required_capabilities().to_vec(),
1065 ),
1066 )
1067 })
1068 .collect();
1069 for binding in plan.capability_bindings() {
1070 dependencies
1071 .get_mut(binding.consumer_instance())
1072 .expect("every resolved binding consumer has Plugin dependencies")
1073 .bindings
1074 .push(PluginDependency::new(
1075 binding.requirement_id(),
1076 binding.capability_id(),
1077 binding.provider_instance(),
1078 binding.provider_order(),
1079 endpoints
1080 .iter()
1081 .find(|((consumer, capability), _)| {
1082 consumer == binding.consumer_instance()
1083 && *capability == binding.capability_id()
1084 })
1085 .and_then(|(_, endpoints)| {
1086 endpoints.iter().find(|endpoint| {
1087 endpoint.requirement_id == binding.requirement_id()
1088 && endpoint.plugin_instance == binding.provider_instance()
1089 })
1090 })
1091 .map(|endpoint| PluginDependencyHandle {
1092 binding: endpoint.clone(),
1093 caller_instance: binding.consumer_instance().to_owned(),
1094 runtime: runtime.clone(),
1095 }),
1096 stream_endpoints
1097 .iter()
1098 .find(|((consumer, capability), _)| {
1099 consumer == binding.consumer_instance()
1100 && *capability == binding.capability_id()
1101 })
1102 .and_then(|(_, endpoints)| {
1103 endpoints.iter().find(|endpoint| {
1104 endpoint.requirement_id == binding.requirement_id()
1105 && endpoint.plugin_instance == binding.provider_instance()
1106 })
1107 })
1108 .map(|endpoint| PluginStreamDependencyHandle {
1109 binding: endpoint.clone(),
1110 caller_instance: binding.consumer_instance().to_owned(),
1111 runtime: runtime.clone(),
1112 }),
1113 event_endpoints
1114 .iter()
1115 .find(|((consumer, capability), _)| {
1116 consumer == binding.consumer_instance()
1117 && *capability == binding.capability_id()
1118 })
1119 .and_then(|(_, endpoints)| {
1120 endpoints.iter().find(|endpoint| {
1121 endpoint.requirement_id == binding.requirement_id()
1122 && endpoint.plugin_instance == binding.provider_instance()
1123 })
1124 })
1125 .map(|endpoint| PluginEventDependencyHandle {
1126 binding: endpoint.clone(),
1127 caller_instance: binding.consumer_instance().to_owned(),
1128 runtime: runtime.clone(),
1129 }),
1130 ));
1131 }
1132 dependencies
1133}