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