1use platform_admin_data::{
22 AdminModule, AdminModuleMetadata, AdminModuleSourceDiagnostics, AdminRemoteModuleDiagnostics,
23};
24use platform_core::error::ErrorDetail;
25use platform_core::{
26 ActorContext, AppContext, AppError, CorrelationId, ErrorCode, EventHandlerRegistry, Migration,
27 PLATFORM_MIGRATIONS, RuntimeConfigDescriptor, RuntimeConfigGroupDescriptor, RuntimeConfigScope,
28 RuntimeConfigType, StoryDisplayDescriptor, StoryDisplaySource, TraceContext,
29};
30use platform_http::ApiOpenApiRouter;
31use platform_module::{
32 AdminSchema, AdminSurface, EventHandlerRegistrationContext, LifecycleActivationRunPolicy,
33 LifecycleStartupCheckKind, LinkedBinding, Module, ModuleHttpMethod, ModuleLoadStatus,
34 ModuleManifest, ModuleSource,
35};
36use platform_module_remote::{RemoteHttpProxyRegistry, RemoteModuleConfig, RemoteModuleSource};
37use platform_runtime::{
38 EnqueueFunctionRequest, FunctionRegistry, RUNTIME_MIGRATIONS, RuntimeClient,
39};
40use std::fs::{self, OpenOptions};
41use std::io::Write as _;
42use std::path::{Path, PathBuf};
43use std::process::{Child, Command};
44use std::sync::Arc;
45use std::thread;
46use std::time::{Duration, Instant};
47
48const DEFAULT_MODULE_SERVICES_FILE: &str = ".lenso/module-services.json";
49const DEFAULT_REMOTE_SERVICE_READY_TIMEOUT_MS: u64 = 10_000;
50const REMOTE_SERVICE_TERMINATE_GRACE_MS: u64 = 800;
51const AUTH_SESSION_CACHE_MAX_TTL: Duration = Duration::from_secs(12 * 60 * 60);
52
53struct LinkedModuleEntry {
54 module_name: &'static str,
55 manifest: fn() -> ModuleManifest,
56 load: fn(&AppContext) -> Module,
57 http_binding: Option<fn() -> LinkedBinding>,
58}
59
60const MODULES_CONFIG_GROUP: RuntimeConfigGroupDescriptor = RuntimeConfigGroupDescriptor {
61 id: "modules",
62 label: "Modules",
63 description: "Module load toggles applied on service startup.",
64 order: 10,
65};
66
67#[derive(Debug, Clone, Copy)]
68pub struct HostLinkedModule {
69 pub module_name: &'static str,
70 pub manifest: fn() -> ModuleManifest,
71 pub load: Option<fn(&AppContext) -> Module>,
72 pub http_binding: Option<fn() -> LinkedBinding>,
73 pub migrations: &'static [Migration],
74}
75
76impl HostLinkedModule {
77 #[must_use]
78 pub const fn manifest_only(
79 module_name: &'static str,
80 manifest: fn() -> ModuleManifest,
81 migrations: &'static [Migration],
82 ) -> Self {
83 Self {
84 module_name,
85 manifest,
86 load: None,
87 http_binding: None,
88 migrations,
89 }
90 }
91
92 #[must_use]
93 pub const fn linked(
94 module_name: &'static str,
95 manifest: fn() -> ModuleManifest,
96 load: fn(&AppContext) -> Module,
97 migrations: &'static [Migration],
98 ) -> Self {
99 Self {
100 module_name,
101 manifest,
102 load: Some(load),
103 http_binding: None,
104 migrations,
105 }
106 }
107
108 #[must_use]
109 pub const fn with_http_binding(mut self, http_binding: fn() -> LinkedBinding) -> Self {
110 self.http_binding = Some(http_binding);
111 self
112 }
113}
114
115#[derive(Debug, Clone, Default)]
116pub struct HostComposition {
117 linked_modules: Vec<HostLinkedModule>,
118}
119
120impl HostComposition {
121 #[must_use]
122 pub fn new() -> Self {
123 Self::default()
124 }
125
126 #[must_use]
127 pub fn with_linked_module(mut self, module: HostLinkedModule) -> Self {
128 self.add_linked_module(module);
129 self
130 }
131
132 pub fn add_linked_module(&mut self, module: HostLinkedModule) {
133 self.linked_modules.push(module);
134 }
135
136 #[must_use]
137 pub fn linked_modules(&self) -> &[HostLinkedModule] {
138 &self.linked_modules
139 }
140}
141
142#[derive(Debug, Clone, Copy, PartialEq, Eq)]
143pub enum CompositionProfile {
144 Core,
145 Demo,
146}
147
148impl CompositionProfile {
149 pub fn parse(value: &str) -> platform_core::AppResult<Self> {
150 match value.trim().to_ascii_lowercase().as_str() {
151 "core" => Ok(Self::Core),
152 "demo" => Ok(Self::Demo),
153 other => Err(AppError::validation(
154 "Invalid Lenso composition profile",
155 vec![ErrorDetail {
156 field: Some("module_sources.linked_profile".to_owned()),
157 reason: format!("expected `core` or `demo`, got `{other}`"),
158 }],
159 )),
160 }
161 }
162
163 pub fn from_config(config: &platform_core::AppConfig) -> platform_core::AppResult<Self> {
164 Self::parse(&config.module_sources.linked_profile)
165 }
166}
167
168impl Default for CompositionProfile {
169 fn default() -> Self {
170 Self::Demo
171 }
172}
173
174const CORE_LINKED_MODULE_ENTRIES: &[LinkedModuleEntry] = &[LinkedModuleEntry {
175 module_name: "platform-story",
176 manifest: story::module::manifest,
177 load: story::module::module,
178 http_binding: Some(story::module::binding),
179}];
180
181const DEMO_LINKED_MODULE_ENTRIES: &[LinkedModuleEntry] = &[
182 LinkedModuleEntry {
183 module_name: "auth",
184 manifest: auth::module::manifest,
185 load: auth::module::module,
186 http_binding: Some(auth::module::binding),
187 },
188 LinkedModuleEntry {
189 module_name: "auth-password",
190 manifest: auth_password::module::manifest,
191 load: auth_password::module::module,
192 http_binding: Some(auth_password::module::binding),
193 },
194 LinkedModuleEntry {
195 module_name: "platform-story",
196 manifest: story::module::manifest,
197 load: story::module::module,
198 http_binding: Some(story::module::binding),
199 },
200];
201
202fn linked_module_entries(profile: CompositionProfile) -> &'static [LinkedModuleEntry] {
203 match profile {
204 CompositionProfile::Core => CORE_LINKED_MODULE_ENTRIES,
205 CompositionProfile::Demo => DEMO_LINKED_MODULE_ENTRIES,
206 }
207}
208
209#[must_use]
210pub const fn auth_linked_module() -> HostLinkedModule {
211 HostLinkedModule::linked(
212 auth::module::MODULE_NAME,
213 auth::module::manifest,
214 auth::module::module,
215 auth::migrations::AUTH_MIGRATIONS,
216 )
217 .with_http_binding(auth::module::binding)
218}
219
220#[must_use]
221pub const fn auth_password_linked_module() -> HostLinkedModule {
222 HostLinkedModule::linked(
223 auth_password::module::MODULE_NAME,
224 auth_password::module::manifest,
225 auth_password::module::module,
226 auth_password::migrations::AUTH_PASSWORD_MIGRATIONS,
227 )
228 .with_http_binding(auth_password::module::binding)
229}
230
231fn linked_module_enabled_from_config(config: &platform_core::AppConfig, module_name: &str) -> bool {
232 config
233 .modules
234 .get(module_name)
235 .is_none_or(platform_core::ModuleConfig::is_enabled)
236}
237
238fn module_enabled_config_key(module_name: &str) -> String {
239 format!("modules.{module_name}.enabled")
240}
241
242fn linked_module_enabled(ctx: &AppContext, module_name: &str) -> bool {
243 ctx.runtime_config
244 .snapshot()
245 .raw(&module_enabled_config_key(module_name))
246 .and_then(serde_json::Value::as_bool)
247 .unwrap_or_else(|| linked_module_enabled_from_config(&ctx.config, module_name))
248}
249
250fn first_disabled_dependency(ctx: &AppContext, manifest: fn() -> ModuleManifest) -> Option<String> {
251 (manifest)()
252 .dependencies
253 .into_iter()
254 .find(|dependency| !linked_module_enabled(ctx, dependency))
255}
256
257fn first_disabled_dependency_from_config(
258 config: &platform_core::AppConfig,
259 manifest: fn() -> ModuleManifest,
260) -> Option<String> {
261 (manifest)()
262 .dependencies
263 .into_iter()
264 .find(|dependency| !linked_module_enabled_from_config(config, dependency))
265}
266
267fn linked_module_with_dependencies_enabled(
268 ctx: &AppContext,
269 module_name: &str,
270 manifest: fn() -> ModuleManifest,
271) -> bool {
272 linked_module_enabled(ctx, module_name) && first_disabled_dependency(ctx, manifest).is_none()
273}
274
275fn linked_module_with_dependencies_enabled_from_config(
276 config: &platform_core::AppConfig,
277 module_name: &str,
278 manifest: fn() -> ModuleManifest,
279) -> bool {
280 linked_module_enabled_from_config(config, module_name)
281 && first_disabled_dependency_from_config(config, manifest).is_none()
282}
283
284fn linked_module_disabled_reason(
285 ctx: &AppContext,
286 module_name: &str,
287 manifest: fn() -> ModuleManifest,
288) -> Option<String> {
289 if !linked_module_enabled(ctx, module_name) {
290 return Some("module disabled by configuration".to_owned());
291 }
292 if let Some(dependency) = first_disabled_dependency(ctx, manifest) {
293 return Some(format!("module dependency disabled: {dependency}"));
294 }
295 None
296}
297
298fn remote_module_enabled_from_config(config: &platform_core::AppConfig, module_name: &str) -> bool {
299 config
300 .modules
301 .get(module_name)
302 .is_none_or(platform_core::ModuleConfig::is_enabled)
303}
304
305fn remote_module_enabled(ctx: &AppContext, module_name: &str) -> bool {
306 ctx.runtime_config
307 .snapshot()
308 .raw(&module_enabled_config_key(module_name))
309 .and_then(serde_json::Value::as_bool)
310 .unwrap_or_else(|| remote_module_enabled_from_config(&ctx.config, module_name))
311}
312
313pub fn auth_actor_resolver_for_context(
314 ctx: &AppContext,
315) -> platform_core::AppResult<Option<Arc<dyn platform_core::ActorResolver>>> {
316 auth_actor_resolver_for_context_with_composition(ctx, &HostComposition::default())
317}
318
319pub fn auth_actor_resolver_for_context_with_composition(
320 ctx: &AppContext,
321 composition: &HostComposition,
322) -> platform_core::AppResult<Option<Arc<dyn platform_core::ActorResolver>>> {
323 let profile = CompositionProfile::from_config(&ctx.config)?;
324 let auth_in_profile = linked_module_entries(profile)
325 .iter()
326 .any(|entry| entry.module_name == auth::module::MODULE_NAME);
327 let auth_in_composition = composition
328 .linked_modules()
329 .iter()
330 .any(|entry| entry.module_name == auth::module::MODULE_NAME);
331 if (!auth_in_profile && !auth_in_composition)
332 || !linked_module_enabled(ctx, auth::module::MODULE_NAME)
333 {
334 return Ok(None);
335 }
336
337 let auth_resolver: Arc<dyn platform_core::ActorResolver> =
338 Arc::new(auth::resolver::AuthActorResolver::new_with_session_cache(
339 ctx.db.clone(),
340 ctx.actor_resolver.clone(),
341 auth_session_cache(ctx)?,
342 ));
343
344 let auth_password_enabled = linked_module_with_dependencies_enabled(
345 ctx,
346 auth_password::module::MODULE_NAME,
347 auth_password::module::manifest,
348 );
349 if auth_password_enabled {
350 if let Some(jwt_resolver) =
351 auth_password::module::jwt_actor_resolver(ctx, auth_resolver.clone())?
352 {
353 return Ok(Some(jwt_resolver));
354 }
355 }
356
357 Ok(Some(auth_resolver))
358}
359
360fn auth_session_cache(
361 ctx: &AppContext,
362) -> platform_core::AppResult<Option<Arc<dyn auth::resolver::SessionCache>>> {
363 match auth::config::AuthRuntimeConfig::from_context(ctx).session_cache {
364 auth::config::SessionCacheMode::Database => Ok(None),
365 auth::config::SessionCacheMode::Redis => {
366 let Some(redis) = ctx.redis.clone() else {
367 return Err(AppError::validation(
368 "Redis auth session cache is not configured",
369 vec![ErrorDetail {
370 field: Some("auth.session_cache".to_owned()),
371 reason: "set REDIS_URL when auth.session_cache is redis".to_owned(),
372 }],
373 ));
374 };
375 Ok(Some(Arc::new(auth::redis_cache::RedisSessionCache::new(
376 redis,
377 AUTH_SESSION_CACHE_MAX_TTL,
378 ))))
379 }
380 }
381}
382
383fn linked_module_entries_for_context(
384 ctx: &AppContext,
385) -> platform_core::AppResult<Vec<&'static LinkedModuleEntry>> {
386 Ok(
387 linked_module_entries(CompositionProfile::from_config(&ctx.config)?)
388 .iter()
389 .filter(|entry| {
390 linked_module_with_dependencies_enabled(ctx, entry.module_name, entry.manifest)
391 })
392 .collect(),
393 )
394}
395
396fn linked_module_entries_for_config(
397 config: &platform_core::AppConfig,
398) -> platform_core::AppResult<Vec<&'static LinkedModuleEntry>> {
399 Ok(
400 linked_module_entries(CompositionProfile::from_config(config)?)
401 .iter()
402 .filter(|entry| {
403 linked_module_with_dependencies_enabled_from_config(
404 config,
405 entry.module_name,
406 entry.manifest,
407 )
408 })
409 .collect(),
410 )
411}
412
413fn disabled_linked_module_entries_for_context(
414 ctx: &AppContext,
415) -> platform_core::AppResult<Vec<&'static LinkedModuleEntry>> {
416 Ok(
417 linked_module_entries(CompositionProfile::from_config(&ctx.config)?)
418 .iter()
419 .filter(|entry| {
420 linked_module_disabled_reason(ctx, entry.module_name, entry.manifest).is_some()
421 })
422 .collect(),
423 )
424}
425
426fn host_linked_modules_for_config(
427 config: &platform_core::AppConfig,
428 composition: &HostComposition,
429) -> Vec<HostLinkedModule> {
430 composition
431 .linked_modules()
432 .iter()
433 .copied()
434 .filter(|entry| {
435 linked_module_with_dependencies_enabled_from_config(
436 config,
437 entry.module_name,
438 entry.manifest,
439 )
440 })
441 .collect()
442}
443
444fn host_linked_modules_for_context(
445 ctx: &AppContext,
446 composition: &HostComposition,
447) -> Vec<HostLinkedModule> {
448 composition
449 .linked_modules()
450 .iter()
451 .copied()
452 .filter(|entry| {
453 linked_module_with_dependencies_enabled(ctx, entry.module_name, entry.manifest)
454 })
455 .collect()
456}
457
458fn disabled_host_linked_modules_for_context(
459 ctx: &AppContext,
460 composition: &HostComposition,
461) -> Vec<HostLinkedModule> {
462 composition
463 .linked_modules()
464 .iter()
465 .copied()
466 .filter(|entry| {
467 linked_module_disabled_reason(ctx, entry.module_name, entry.manifest).is_some()
468 })
469 .collect()
470}
471
472fn load_host_linked_module(ctx: &AppContext, entry: HostLinkedModule) -> Module {
473 match entry.load {
474 Some(load) => load(ctx),
475 None => Module::linked((entry.manifest)(), LinkedBinding::builder().build()),
476 }
477}
478
479#[must_use]
484pub fn modules(ctx: &AppContext) -> Vec<Module> {
485 modules_for_profile(ctx, CompositionProfile::default())
486}
487
488pub fn modules_for_config(ctx: &AppContext) -> platform_core::AppResult<Vec<Module>> {
489 Ok(linked_module_entries_for_context(ctx)?
490 .into_iter()
491 .map(|entry| (entry.load)(ctx))
492 .collect())
493}
494
495pub fn modules_for_config_with_composition(
496 ctx: &AppContext,
497 composition: &HostComposition,
498) -> platform_core::AppResult<Vec<Module>> {
499 let mut modules = modules_for_config(ctx)?;
500 modules.extend(
501 host_linked_modules_for_context(ctx, composition)
502 .into_iter()
503 .map(|entry| load_host_linked_module(ctx, entry)),
504 );
505 Ok(modules)
506}
507
508#[must_use]
509pub fn modules_for_profile(ctx: &AppContext, profile: CompositionProfile) -> Vec<Module> {
510 linked_module_entries(profile)
511 .iter()
512 .map(|entry| (entry.load)(ctx))
513 .collect()
514}
515
516pub async fn load_modules(ctx: &AppContext) -> platform_core::AppResult<Vec<Module>> {
522 load_modules_with_composition(ctx, &HostComposition::default()).await
523}
524
525pub async fn load_modules_with_composition(
526 ctx: &AppContext,
527 composition: &HostComposition,
528) -> platform_core::AppResult<Vec<Module>> {
529 let mut loaded = modules_for_config_with_composition(ctx, composition)?;
530
531 for remote in &ctx.config.module_sources.remote {
532 if !remote_module_enabled(ctx, &remote.name) {
533 continue;
534 }
535 let source = RemoteModuleSource::new(remote_module_config(remote))?;
536 loaded.push(source.load().await?);
537 }
538
539 Ok(loaded)
540}
541
542pub fn migrations_for_config(
543 config: &platform_core::AppConfig,
544) -> platform_core::AppResult<Vec<Migration>> {
545 migrations_for_config_with_composition(config, &HostComposition::default())
546}
547
548pub fn migrations_for_config_with_composition(
549 config: &platform_core::AppConfig,
550 composition: &HostComposition,
551) -> platform_core::AppResult<Vec<Migration>> {
552 let mut migrations = PLATFORM_MIGRATIONS
553 .iter()
554 .chain(RUNTIME_MIGRATIONS)
555 .copied()
556 .collect::<Vec<_>>();
557
558 if CompositionProfile::from_config(config)? == CompositionProfile::Demo {
559 if linked_module_enabled_from_config(config, "auth") {
560 migrations.extend(auth::migrations::AUTH_MIGRATIONS.iter().copied());
561 }
562 if linked_module_with_dependencies_enabled_from_config(
563 config,
564 "auth-password",
565 auth_password::module::manifest,
566 ) {
567 migrations.extend(
568 auth_password::migrations::AUTH_PASSWORD_MIGRATIONS
569 .iter()
570 .copied(),
571 );
572 }
573 }
574
575 for module in host_linked_modules_for_config(config, composition) {
576 migrations.extend(module.migrations.iter().copied());
577 }
578
579 Ok(migrations)
580}
581
582#[must_use]
583pub fn migrations_for_profile(profile: CompositionProfile) -> Vec<Migration> {
584 let mut migrations = PLATFORM_MIGRATIONS
585 .iter()
586 .chain(RUNTIME_MIGRATIONS)
587 .copied()
588 .collect::<Vec<_>>();
589
590 if profile == CompositionProfile::Demo {
591 migrations.extend(auth::migrations::AUTH_MIGRATIONS.iter().copied());
592 migrations.extend(
593 auth_password::migrations::AUTH_PASSWORD_MIGRATIONS
594 .iter()
595 .copied(),
596 );
597 }
598
599 migrations
600}
601
602#[must_use]
605pub fn module_manifests() -> Vec<ModuleManifest> {
606 module_manifests_for_profile(CompositionProfile::default())
607}
608
609#[must_use]
610pub fn module_manifests_for_profile(profile: CompositionProfile) -> Vec<ModuleManifest> {
611 linked_module_entries(profile)
612 .iter()
613 .map(|entry| (entry.manifest)())
614 .collect()
615}
616
617#[must_use]
619pub fn linked_runtime_function_declaration_sources() -> Vec<(
620 String,
621 ModuleSource,
622 Option<platform_module::RuntimeSurface>,
623)> {
624 linked_runtime_function_declaration_sources_for_profile(CompositionProfile::default())
625}
626
627#[must_use]
628pub fn linked_runtime_function_declaration_sources_for_profile(
629 profile: CompositionProfile,
630) -> Vec<(
631 String,
632 ModuleSource,
633 Option<platform_module::RuntimeSurface>,
634)> {
635 module_manifests_for_profile(profile)
636 .into_iter()
637 .map(|manifest| (manifest.name, ModuleSource::Linked, manifest.runtime))
638 .collect()
639}
640
641pub fn linked_runtime_function_declaration_sources_for_config(
642 config: &platform_core::AppConfig,
643) -> platform_core::AppResult<
644 Vec<(
645 String,
646 ModuleSource,
647 Option<platform_module::RuntimeSurface>,
648 )>,
649> {
650 Ok(linked_module_entries_for_config(config)?
651 .into_iter()
652 .map(|entry| {
653 let manifest = (entry.manifest)();
654 (manifest.name, ModuleSource::Linked, manifest.runtime)
655 })
656 .collect())
657}
658
659pub fn linked_runtime_function_declaration_sources_for_context(
660 ctx: &AppContext,
661) -> platform_core::AppResult<
662 Vec<(
663 String,
664 ModuleSource,
665 Option<platform_module::RuntimeSurface>,
666 )>,
667> {
668 Ok(linked_module_entries_for_context(ctx)?
669 .into_iter()
670 .map(|entry| {
671 let manifest = (entry.manifest)();
672 (manifest.name, ModuleSource::Linked, manifest.runtime)
673 })
674 .collect())
675}
676
677pub fn linked_runtime_function_declaration_sources_for_context_with_composition(
678 ctx: &AppContext,
679 composition: &HostComposition,
680) -> platform_core::AppResult<
681 Vec<(
682 String,
683 ModuleSource,
684 Option<platform_module::RuntimeSurface>,
685 )>,
686> {
687 let mut sources = linked_runtime_function_declaration_sources_for_context(ctx)?;
688 sources.extend(
689 host_linked_modules_for_context(ctx, composition)
690 .into_iter()
691 .map(|entry| {
692 let manifest = (entry.manifest)();
693 (manifest.name, ModuleSource::Linked, manifest.runtime)
694 }),
695 );
696 Ok(sources)
697}
698
699#[must_use]
702pub fn runtime_function_declaration_sources_from_metadata(
703 modules: &[AdminModuleMetadata],
704) -> Vec<(
705 String,
706 ModuleSource,
707 Option<platform_module::RuntimeSurface>,
708)> {
709 modules
710 .iter()
711 .filter(|module| matches!(module.load_status, ModuleLoadStatus::Loaded))
712 .map(|module| {
713 (
714 module.module_name.clone(),
715 module.source,
716 module.runtime.clone(),
717 )
718 })
719 .collect()
720}
721
722#[derive(Debug, Clone, PartialEq, Eq)]
727pub struct LinkedHttpRouteOwner {
728 pub module_name: String,
729 pub public_prefixes: &'static [&'static str],
730}
731
732#[must_use]
733pub fn linked_http_route_owners() -> Vec<LinkedHttpRouteOwner> {
734 linked_http_route_owners_for_profile(CompositionProfile::default())
735}
736
737#[must_use]
738pub fn linked_http_route_owners_for_profile(
739 profile: CompositionProfile,
740) -> Vec<LinkedHttpRouteOwner> {
741 linked_module_entries(profile)
742 .iter()
743 .filter_map(|entry| {
744 let http = entry.http_binding?().http?;
745 Some(LinkedHttpRouteOwner {
746 module_name: entry.module_name.to_owned(),
747 public_prefixes: http.public_prefixes,
748 })
749 })
750 .collect()
751}
752
753#[must_use]
755pub fn linked_http_modules() -> Vec<Module> {
756 linked_http_modules_for_profile(CompositionProfile::default())
757}
758
759#[must_use]
760pub fn linked_http_modules_for_profile(profile: CompositionProfile) -> Vec<Module> {
761 linked_module_entries(profile)
762 .iter()
763 .filter_map(|entry| {
764 let http_binding = entry.http_binding?;
765 Some(Module::linked((entry.manifest)(), http_binding()))
766 })
767 .collect()
768}
769
770pub fn linked_http_modules_for_config(
771 config: &platform_core::AppConfig,
772) -> platform_core::AppResult<Vec<Module>> {
773 Ok(linked_module_entries_for_config(config)?
774 .into_iter()
775 .filter_map(|entry| {
776 let http_binding = entry.http_binding?;
777 Some(Module::linked((entry.manifest)(), http_binding()))
778 })
779 .collect())
780}
781
782pub fn linked_http_modules_for_context(ctx: &AppContext) -> platform_core::AppResult<Vec<Module>> {
783 Ok(linked_module_entries_for_context(ctx)?
784 .into_iter()
785 .filter_map(|entry| {
786 let http_binding = entry.http_binding?;
787 Some(Module::linked((entry.manifest)(), http_binding()))
788 })
789 .collect())
790}
791
792pub fn linked_http_modules_for_context_with_composition(
793 ctx: &AppContext,
794 composition: &HostComposition,
795) -> platform_core::AppResult<Vec<Module>> {
796 let mut modules = linked_http_modules_for_context(ctx)?;
797 modules.extend(
798 host_linked_modules_for_context(ctx, composition)
799 .into_iter()
800 .filter_map(|entry| {
801 let http_binding = entry.http_binding?;
802 Some(Module::linked((entry.manifest)(), http_binding()))
803 }),
804 );
805 Ok(modules)
806}
807
808#[must_use]
813pub fn admin_modules(ctx: &AppContext) -> Vec<AdminModule> {
814 admin_modules_from_modules(modules(ctx))
815}
816
817pub async fn load_admin_modules(ctx: &AppContext) -> platform_core::AppResult<Vec<AdminModule>> {
819 load_admin_modules_with_composition(ctx, &HostComposition::default()).await
820}
821
822pub async fn load_admin_modules_with_composition(
823 ctx: &AppContext,
824 composition: &HostComposition,
825) -> platform_core::AppResult<Vec<AdminModule>> {
826 let mut admin_modules =
827 admin_modules_from_modules(modules_for_config_with_composition(ctx, composition)?);
828
829 for remote in &ctx.config.module_sources.remote {
830 if !remote_module_enabled(ctx, &remote.name) {
831 continue;
832 }
833 let source = RemoteModuleSource::new(remote_module_config(remote))?;
834 match source.load().await {
835 Ok(module) => admin_modules.extend(admin_modules_from_modules(vec![module])),
836 Err(error) => admin_modules.push(failed_remote_admin_module(
837 remote.name.clone(),
838 error.public_message,
839 )),
840 }
841 }
842
843 Ok(admin_modules)
844}
845
846pub async fn load_admin_module_metadata(
850 ctx: &AppContext,
851) -> platform_core::AppResult<Vec<AdminModuleMetadata>> {
852 load_admin_module_metadata_with_composition(ctx, &HostComposition::default()).await
853}
854
855pub async fn load_admin_module_metadata_with_composition(
856 ctx: &AppContext,
857 composition: &HostComposition,
858) -> platform_core::AppResult<Vec<AdminModuleMetadata>> {
859 let mut metadata =
860 admin_metadata_from_modules(modules_for_config_with_composition(ctx, composition)?);
861 metadata.extend(disabled_linked_admin_metadata(ctx)?);
862 metadata.extend(disabled_host_linked_admin_metadata(ctx, composition));
863
864 for remote in &ctx.config.module_sources.remote {
865 let config = remote_module_config(remote);
866 if !remote_module_enabled(ctx, &remote.name) {
867 metadata.push(disabled_remote_admin_metadata(&config));
868 continue;
869 }
870 let checked_at = current_timestamp();
871 let source = RemoteModuleSource::new(config.clone())?;
872 let load_started = Instant::now();
873 match source.load().await {
874 Ok(module) => metadata.extend(remote_admin_metadata_from_module(
875 module,
876 &config,
877 checked_at,
878 Some(duration_ms(load_started)),
879 None,
880 )),
881 Err(error) => metadata.push(failed_remote_admin_metadata(
882 &config,
883 Some(checked_at),
884 Some(duration_ms(load_started)),
885 error.public_message,
886 )),
887 }
888 }
889
890 Ok(metadata)
891}
892
893pub async fn load_remote_http_proxy_registry(
894 ctx: &AppContext,
895) -> platform_core::AppResult<RemoteHttpProxyRegistry> {
896 let mut remote_modules = Vec::new();
897 let mut remote_configs = Vec::new();
898
899 for remote in &ctx.config.module_sources.remote {
900 if !remote_module_enabled(ctx, &remote.name) {
901 continue;
902 }
903 let config = remote_module_config(remote);
904 let source = RemoteModuleSource::new(config.clone())?;
905 if let Ok(module) = source.load().await {
906 remote_modules.push(module);
907 remote_configs.push(config);
908 }
909 }
910
911 Ok(RemoteHttpProxyRegistry::from_modules(
912 &remote_modules,
913 &remote_configs,
914 ))
915}
916
917fn admin_modules_from_modules(modules: Vec<Module>) -> Vec<AdminModule> {
918 modules
919 .into_iter()
920 .filter_map(|module| {
921 let data_source = module.admin_data;
923 let action_source = module.admin_actions;
924 let query_source = module.admin_queries;
925 if data_source.is_none() && action_source.is_none() && query_source.is_none() {
926 return None;
927 }
928 let ModuleManifest { name, admin, .. } = module.manifest;
929 let admin = admin?;
930 let (schema, listed_in_schema) = match &admin {
931 AdminSurface::Schema(schema) => (schema.clone(), true),
932 AdminSurface::DeclarativeCustom(surface) => (
933 surface.fallback_schema.clone().unwrap_or(AdminSchema {
934 entities: Vec::new(),
935 }),
936 false,
937 ),
938 AdminSurface::EmbeddedCustom(_) => return None,
939 _ => return None,
940 };
941 Some(AdminModule {
942 module_name: name,
943 source: module.source,
944 load_status: module.load_status,
945 schema,
946 admin: Some(admin),
947 listed_in_schema,
948 data_source,
949 action_source,
950 query_source,
951 })
952 })
953 .collect()
954}
955
956fn admin_metadata_from_modules(modules: Vec<Module>) -> Vec<AdminModuleMetadata> {
957 modules
958 .into_iter()
959 .map(|module| {
960 let ModuleManifest {
961 name,
962 admin,
963 http_routes,
964 runtime,
965 events,
966 lifecycle,
967 console,
968 story_display,
969 capabilities,
970 dependencies,
971 ..
972 } = module.manifest;
973 AdminModuleMetadata {
974 module_name: name,
975 source: module.source,
976 load_status: module.load_status,
977 http_routes,
978 runtime,
979 events,
980 lifecycle,
981 console,
982 story_display,
983 capabilities,
984 dependencies,
985 admin,
986 source_diagnostics: None,
987 }
988 })
989 .collect()
990}
991
992fn failed_remote_admin_module(name: String, message: String) -> AdminModule {
993 AdminModule {
994 module_name: name,
995 source: ModuleSource::Remote,
996 load_status: ModuleLoadStatus::Error { message },
997 schema: AdminSchema {
998 entities: Vec::new(),
999 },
1000 admin: None,
1001 listed_in_schema: true,
1002 data_source: None,
1003 action_source: None,
1004 query_source: None,
1005 }
1006}
1007
1008fn remote_admin_metadata_from_module(
1009 module: Module,
1010 config: &RemoteModuleConfig,
1011 checked_at: String,
1012 load_duration_ms: Option<u64>,
1013 load_error: Option<String>,
1014) -> Vec<AdminModuleMetadata> {
1015 admin_metadata_from_modules(vec![module])
1016 .into_iter()
1017 .map(|mut metadata| {
1018 metadata.source_diagnostics = Some(remote_source_diagnostics(
1019 config,
1020 Some(checked_at.clone()),
1021 load_duration_ms,
1022 load_error.clone(),
1023 ));
1024 metadata
1025 })
1026 .collect()
1027}
1028
1029fn failed_remote_admin_metadata(
1030 config: &RemoteModuleConfig,
1031 checked_at: Option<String>,
1032 load_duration_ms: Option<u64>,
1033 message: String,
1034) -> AdminModuleMetadata {
1035 AdminModuleMetadata {
1036 module_name: config.name.clone(),
1037 source: ModuleSource::Remote,
1038 load_status: ModuleLoadStatus::Error {
1039 message: message.clone(),
1040 },
1041 http_routes: Vec::new(),
1042 runtime: None,
1043 events: None,
1044 lifecycle: None,
1045 console: Vec::new(),
1046 story_display: Vec::new(),
1047 capabilities: Vec::new(),
1048 dependencies: Vec::new(),
1049 admin: None,
1050 source_diagnostics: Some(remote_source_diagnostics(
1051 config,
1052 checked_at,
1053 load_duration_ms,
1054 Some(message),
1055 )),
1056 }
1057}
1058
1059fn disabled_remote_admin_metadata(config: &RemoteModuleConfig) -> AdminModuleMetadata {
1060 AdminModuleMetadata {
1061 module_name: config.name.clone(),
1062 source: ModuleSource::Remote,
1063 load_status: ModuleLoadStatus::Error {
1064 message: "module disabled by configuration".to_owned(),
1065 },
1066 http_routes: Vec::new(),
1067 runtime: None,
1068 events: None,
1069 lifecycle: None,
1070 console: Vec::new(),
1071 story_display: Vec::new(),
1072 capabilities: Vec::new(),
1073 dependencies: Vec::new(),
1074 admin: None,
1075 source_diagnostics: Some(remote_source_diagnostics(config, None, None, None)),
1076 }
1077}
1078
1079fn disabled_linked_admin_metadata(
1080 ctx: &AppContext,
1081) -> platform_core::AppResult<Vec<AdminModuleMetadata>> {
1082 Ok(disabled_linked_module_entries_for_context(ctx)?
1083 .into_iter()
1084 .map(|entry| {
1085 let ModuleManifest {
1086 name,
1087 admin,
1088 http_routes,
1089 runtime,
1090 events,
1091 lifecycle,
1092 console,
1093 story_display,
1094 capabilities,
1095 dependencies,
1096 ..
1097 } = (entry.manifest)();
1098 AdminModuleMetadata {
1099 module_name: name,
1100 source: ModuleSource::Linked,
1101 load_status: ModuleLoadStatus::Error {
1102 message: linked_module_disabled_reason(ctx, entry.module_name, entry.manifest)
1103 .unwrap_or_else(|| "module disabled by configuration".to_owned()),
1104 },
1105 http_routes,
1106 runtime,
1107 events,
1108 lifecycle,
1109 console,
1110 story_display,
1111 capabilities,
1112 dependencies,
1113 admin,
1114 source_diagnostics: None,
1115 }
1116 })
1117 .collect())
1118}
1119
1120fn disabled_host_linked_admin_metadata(
1121 ctx: &AppContext,
1122 composition: &HostComposition,
1123) -> Vec<AdminModuleMetadata> {
1124 disabled_host_linked_modules_for_context(ctx, composition)
1125 .into_iter()
1126 .map(|entry| {
1127 let ModuleManifest {
1128 name,
1129 admin,
1130 http_routes,
1131 runtime,
1132 events,
1133 lifecycle,
1134 console,
1135 story_display,
1136 capabilities,
1137 dependencies,
1138 ..
1139 } = (entry.manifest)();
1140 AdminModuleMetadata {
1141 module_name: name,
1142 source: ModuleSource::Linked,
1143 load_status: ModuleLoadStatus::Error {
1144 message: linked_module_disabled_reason(ctx, entry.module_name, entry.manifest)
1145 .unwrap_or_else(|| "module disabled by configuration".to_owned()),
1146 },
1147 http_routes,
1148 runtime,
1149 events,
1150 lifecycle,
1151 console,
1152 story_display,
1153 capabilities,
1154 dependencies,
1155 admin,
1156 source_diagnostics: None,
1157 }
1158 })
1159 .collect()
1160}
1161
1162fn remote_source_diagnostics(
1163 config: &RemoteModuleConfig,
1164 checked_at: Option<String>,
1165 load_duration_ms: Option<u64>,
1166 load_error: Option<String>,
1167) -> AdminModuleSourceDiagnostics {
1168 let (transport, manifest_url) = match config.transport {
1169 platform_module_remote::RemoteModuleTransport::HttpJson => {
1170 ("http_json", format!("{}/manifest", config.base_url))
1171 }
1172 platform_module_remote::RemoteModuleTransport::Grpc => (
1173 "grpc",
1174 format!(
1175 "{}#lenso.remote.v1.RemoteModule/GetManifest",
1176 config.base_url
1177 ),
1178 ),
1179 };
1180 AdminModuleSourceDiagnostics::Remote(AdminRemoteModuleDiagnostics {
1181 transport: transport.to_owned(),
1182 base_url: config.base_url.clone(),
1183 manifest_url,
1184 timeout_ms: config.timeout_ms,
1185 auth_configured: config.auth_token.is_some(),
1186 load_duration_ms,
1187 last_checked_at: checked_at,
1188 last_load_error: load_error,
1189 })
1190}
1191
1192fn remote_module_config(source: &platform_core::RemoteModuleSourceConfig) -> RemoteModuleConfig {
1193 let mut config = RemoteModuleConfig::new(source.name.clone(), source.base_url.clone())
1194 .with_timeout_ms(source.timeout_ms);
1195
1196 if let Some(env_name) = &source.auth_token_env {
1197 if let Ok(token) = std::env::var(env_name) {
1198 config = config.with_auth_token(token);
1199 }
1200 }
1201
1202 config
1203}
1204
1205#[derive(Debug)]
1206pub struct RemoteModuleServiceSupervisor {
1207 services: Vec<RemoteModuleServiceHandle>,
1208}
1209
1210impl RemoteModuleServiceSupervisor {
1211 #[must_use]
1212 pub fn is_empty(&self) -> bool {
1213 self.services.is_empty()
1214 }
1215}
1216
1217impl Drop for RemoteModuleServiceSupervisor {
1218 fn drop(&mut self) {
1219 for service in &mut self.services {
1220 terminate_remote_module_service(&mut service.child);
1221 release_remote_module_service_state(&service.lock_file_path, &service.pid_file_path);
1222 }
1223 }
1224}
1225
1226#[derive(Debug)]
1227struct RemoteModuleServiceHandle {
1228 child: Child,
1229 lock_file_path: PathBuf,
1230 pid_file_path: PathBuf,
1231}
1232
1233#[derive(Debug, Clone, PartialEq, Eq)]
1234struct RemoteModuleServiceSpec {
1235 module_name: String,
1236 service_name: String,
1237 command: String,
1238 cwd: Option<PathBuf>,
1239 ready_url: String,
1240 ready_timeout_ms: u64,
1241 auto_start: bool,
1242}
1243
1244pub async fn start_installed_remote_module_services(
1245 ctx: &AppContext,
1246) -> platform_core::AppResult<RemoteModuleServiceSupervisor> {
1247 start_installed_remote_module_services_from_path(ctx, Path::new(DEFAULT_MODULE_SERVICES_FILE))
1248 .await
1249}
1250
1251pub async fn start_installed_remote_module_services_from_path(
1252 ctx: &AppContext,
1253 services_file_path: &Path,
1254) -> platform_core::AppResult<RemoteModuleServiceSupervisor> {
1255 let specs = read_remote_module_service_specs(services_file_path)?;
1256 let services_state_dir = services_file_path
1257 .parent()
1258 .unwrap_or_else(|| Path::new("."));
1259 let client = reqwest::Client::builder()
1260 .timeout(Duration::from_millis(800))
1261 .build()
1262 .map_err(|source| {
1263 AppError::new(ErrorCode::Internal, "failed to build HTTP client").with_source(source)
1264 })?;
1265 let mut services = Vec::new();
1266
1267 for spec in specs {
1268 if !spec.auto_start || !remote_service_module_enabled(ctx, &spec.module_name) {
1269 continue;
1270 }
1271 if remote_service_ready(&client, &spec.ready_url).await {
1272 tracing::info!(
1273 module = %spec.module_name,
1274 service = %spec.service_name,
1275 ready_url = %spec.ready_url,
1276 "remote module service already ready"
1277 );
1278 continue;
1279 }
1280 let lock_file_path = remote_module_service_state_path(services_state_dir, &spec, "lock");
1281 let pid_file_path = remote_module_service_state_path(services_state_dir, &spec, "pid");
1282 if !claim_remote_module_service_lock(&client, &spec, &lock_file_path, &pid_file_path)
1283 .await?
1284 {
1285 continue;
1286 }
1287 let mut child = match spawn_remote_module_service(&spec) {
1288 Ok(child) => child,
1289 Err(error) => {
1290 release_remote_module_service_state(&lock_file_path, &pid_file_path);
1291 return Err(error);
1292 }
1293 };
1294 if let Err(error) = write_remote_module_service_pid(&pid_file_path, child.id()) {
1295 terminate_remote_module_service(&mut child);
1296 release_remote_module_service_state(&lock_file_path, &pid_file_path);
1297 return Err(error);
1298 }
1299 if let Err(error) = wait_for_remote_module_service(&client, &spec, &mut child).await {
1300 terminate_remote_module_service(&mut child);
1301 release_remote_module_service_state(&lock_file_path, &pid_file_path);
1302 return Err(error);
1303 }
1304 tracing::info!(
1305 module = %spec.module_name,
1306 service = %spec.service_name,
1307 ready_url = %spec.ready_url,
1308 "started remote module service"
1309 );
1310 services.push(RemoteModuleServiceHandle {
1311 child,
1312 lock_file_path,
1313 pid_file_path,
1314 });
1315 }
1316
1317 Ok(RemoteModuleServiceSupervisor { services })
1318}
1319
1320fn remote_service_module_enabled(ctx: &AppContext, module_name: &str) -> bool {
1321 ctx.config
1322 .module_sources
1323 .remote
1324 .iter()
1325 .any(|remote| remote.name == module_name)
1326 && remote_module_enabled(ctx, module_name)
1327}
1328
1329fn read_remote_module_service_specs(
1330 services_file_path: &Path,
1331) -> platform_core::AppResult<Vec<RemoteModuleServiceSpec>> {
1332 let source = match std::fs::read_to_string(services_file_path) {
1333 Ok(source) => source,
1334 Err(error) if error.kind() == std::io::ErrorKind::NotFound => return Ok(Vec::new()),
1335 Err(source) => {
1336 return Err(AppError::new(
1337 ErrorCode::ExternalDependency,
1338 format!("remote module services file could not be read: {source}"),
1339 ));
1340 }
1341 };
1342 let value = serde_json::from_str::<serde_json::Value>(&source).map_err(|source| {
1343 AppError::new(
1344 ErrorCode::Validation,
1345 format!("remote module services file could not be parsed: {source}"),
1346 )
1347 })?;
1348 parse_remote_module_service_specs(&value)
1349}
1350
1351fn parse_remote_module_service_specs(
1352 value: &serde_json::Value,
1353) -> platform_core::AppResult<Vec<RemoteModuleServiceSpec>> {
1354 let modules = value
1355 .get("modules")
1356 .and_then(serde_json::Value::as_array)
1357 .ok_or_else(|| {
1358 AppError::new(
1359 ErrorCode::Validation,
1360 "remote module services file modules must be an array",
1361 )
1362 })?;
1363 let mut specs = Vec::new();
1364 for module in modules {
1365 let module_name = json_string(module, "moduleName")?;
1366 let services = module
1367 .get("services")
1368 .and_then(serde_json::Value::as_array)
1369 .ok_or_else(|| {
1370 AppError::new(
1371 ErrorCode::Validation,
1372 format!("{module_name} services must be an array"),
1373 )
1374 })?;
1375 for service in services {
1376 let command = json_string(service, "command")?;
1377 let ready_url = json_string(service, "readyUrl")?;
1378 specs.push(RemoteModuleServiceSpec {
1379 module_name: module_name.clone(),
1380 service_name: service
1381 .get("name")
1382 .and_then(serde_json::Value::as_str)
1383 .unwrap_or(&module_name)
1384 .to_owned(),
1385 command,
1386 cwd: service
1387 .get("cwd")
1388 .and_then(serde_json::Value::as_str)
1389 .map(PathBuf::from),
1390 ready_url,
1391 ready_timeout_ms: service
1392 .get("readyTimeoutMs")
1393 .and_then(serde_json::Value::as_u64)
1394 .unwrap_or(DEFAULT_REMOTE_SERVICE_READY_TIMEOUT_MS),
1395 auto_start: service
1396 .get("autoStart")
1397 .and_then(serde_json::Value::as_bool)
1398 .unwrap_or(true),
1399 });
1400 }
1401 }
1402 Ok(specs)
1403}
1404
1405fn json_string(value: &serde_json::Value, key: &str) -> platform_core::AppResult<String> {
1406 value
1407 .get(key)
1408 .and_then(serde_json::Value::as_str)
1409 .map(str::to_owned)
1410 .ok_or_else(|| AppError::new(ErrorCode::Validation, format!("{key} must be a string")))
1411}
1412
1413fn spawn_remote_module_service(spec: &RemoteModuleServiceSpec) -> platform_core::AppResult<Child> {
1414 let cwd = spec
1415 .cwd
1416 .clone()
1417 .unwrap_or(std::env::current_dir().map_err(|source| {
1418 AppError::new(ErrorCode::Internal, "failed to resolve current directory")
1419 .with_source(source)
1420 })?);
1421 let mut command = shell_command(&spec.command);
1422 command.current_dir(cwd);
1423 configure_remote_module_service_process(&mut command);
1424 command.spawn().map_err(|source| {
1425 AppError::new(
1426 ErrorCode::ExternalDependency,
1427 format!(
1428 "failed to start remote module service {}: {}",
1429 spec.module_name, spec.service_name
1430 ),
1431 )
1432 .with_source(source)
1433 })
1434}
1435
1436async fn wait_for_remote_module_service(
1437 client: &reqwest::Client,
1438 spec: &RemoteModuleServiceSpec,
1439 child: &mut Child,
1440) -> platform_core::AppResult<()> {
1441 let started = Instant::now();
1442 let timeout = Duration::from_millis(spec.ready_timeout_ms);
1443 loop {
1444 if remote_service_ready(client, &spec.ready_url).await {
1445 return Ok(());
1446 }
1447 if let Some(status) = child.try_wait().map_err(|source| {
1448 AppError::new(
1449 ErrorCode::ExternalDependency,
1450 format!(
1451 "remote module service {} status could not be checked",
1452 spec.service_name
1453 ),
1454 )
1455 .with_source(source)
1456 })? {
1457 return Err(AppError::new(
1458 ErrorCode::ExternalDependency,
1459 format!(
1460 "remote module service {} exited before it became ready: {status}",
1461 spec.service_name
1462 ),
1463 ));
1464 }
1465 if started.elapsed() >= timeout {
1466 return Err(AppError::new(
1467 ErrorCode::ExternalDependency,
1468 format!(
1469 "remote module service {} did not become ready at {}",
1470 spec.service_name, spec.ready_url
1471 ),
1472 ));
1473 }
1474 tokio::time::sleep(Duration::from_millis(200)).await;
1475 }
1476}
1477
1478async fn remote_service_ready(client: &reqwest::Client, ready_url: &str) -> bool {
1479 client
1480 .get(ready_url)
1481 .send()
1482 .await
1483 .is_ok_and(|response| response.status().is_success())
1484}
1485
1486async fn claim_remote_module_service_lock(
1487 client: &reqwest::Client,
1488 spec: &RemoteModuleServiceSpec,
1489 lock_file_path: &Path,
1490 pid_file_path: &Path,
1491) -> platform_core::AppResult<bool> {
1492 match create_remote_module_service_lock(lock_file_path) {
1493 Ok(()) => return Ok(true),
1494 Err(error) if error.kind() == std::io::ErrorKind::AlreadyExists => {}
1495 Err(source) => {
1496 return Err(AppError::new(
1497 ErrorCode::ExternalDependency,
1498 format!(
1499 "remote module service {} lock could not be created: {source}",
1500 spec.service_name
1501 ),
1502 ));
1503 }
1504 }
1505
1506 tracing::info!(
1507 module = %spec.module_name,
1508 service = %spec.service_name,
1509 ready_url = %spec.ready_url,
1510 "remote module service startup already claimed"
1511 );
1512 if wait_for_remote_module_service_ready(
1513 client,
1514 &spec.ready_url,
1515 Duration::from_millis(spec.ready_timeout_ms),
1516 )
1517 .await
1518 {
1519 return Ok(false);
1520 }
1521
1522 tracing::warn!(
1523 module = %spec.module_name,
1524 service = %spec.service_name,
1525 lock_file = %lock_file_path.display(),
1526 "remote module service lock did not become ready before timeout; treating it as stale"
1527 );
1528 terminate_stale_remote_module_service(pid_file_path);
1529 release_remote_module_service_state(lock_file_path, pid_file_path);
1530 match create_remote_module_service_lock(lock_file_path) {
1531 Ok(()) => Ok(true),
1532 Err(error) if error.kind() == std::io::ErrorKind::AlreadyExists => Ok(false),
1533 Err(source) => Err(AppError::new(
1534 ErrorCode::ExternalDependency,
1535 format!(
1536 "stale remote module service {} lock could not be replaced: {source}",
1537 spec.service_name
1538 ),
1539 )),
1540 }
1541}
1542
1543fn create_remote_module_service_lock(lock_file_path: &Path) -> std::io::Result<()> {
1544 if let Some(parent) = lock_file_path.parent() {
1545 fs::create_dir_all(parent)?;
1546 }
1547 let mut file = OpenOptions::new()
1548 .write(true)
1549 .create_new(true)
1550 .open(lock_file_path)?;
1551 writeln!(file, "owner_pid={}", std::process::id())?;
1552 Ok(())
1553}
1554
1555fn write_remote_module_service_pid(
1556 pid_file_path: &Path,
1557 child_pid: u32,
1558) -> platform_core::AppResult<()> {
1559 if let Some(parent) = pid_file_path.parent() {
1560 fs::create_dir_all(parent).map_err(|source| {
1561 AppError::new(
1562 ErrorCode::ExternalDependency,
1563 format!("remote module service pid directory could not be created: {source}"),
1564 )
1565 })?;
1566 }
1567 fs::write(pid_file_path, format!("{child_pid}\n")).map_err(|source| {
1568 AppError::new(
1569 ErrorCode::ExternalDependency,
1570 format!("remote module service pid file could not be written: {source}"),
1571 )
1572 })
1573}
1574
1575fn release_remote_module_service_state(lock_file_path: &Path, pid_file_path: &Path) {
1576 let _ = fs::remove_file(pid_file_path);
1577 let _ = fs::remove_file(lock_file_path);
1578}
1579
1580#[cfg(unix)]
1581fn terminate_stale_remote_module_service(pid_file_path: &Path) {
1582 let Ok(source) = fs::read_to_string(pid_file_path) else {
1583 return;
1584 };
1585 let Ok(pid) = source.trim().parse::<u32>() else {
1586 return;
1587 };
1588
1589 let _ = Command::new("kill")
1590 .arg("-TERM")
1591 .arg(format!("-{pid}"))
1592 .status();
1593 thread::sleep(Duration::from_millis(100));
1594}
1595
1596#[cfg(not(unix))]
1597fn terminate_stale_remote_module_service(_pid_file_path: &Path) {}
1598
1599async fn wait_for_remote_module_service_ready(
1600 client: &reqwest::Client,
1601 ready_url: &str,
1602 timeout: Duration,
1603) -> bool {
1604 let started = Instant::now();
1605 loop {
1606 if remote_service_ready(client, ready_url).await {
1607 return true;
1608 }
1609 if started.elapsed() >= timeout {
1610 return false;
1611 }
1612 tokio::time::sleep(Duration::from_millis(200)).await;
1613 }
1614}
1615
1616fn remote_module_service_state_path(
1617 services_state_dir: &Path,
1618 spec: &RemoteModuleServiceSpec,
1619 extension: &str,
1620) -> PathBuf {
1621 services_state_dir.join(format!(
1622 "remote-{}-{}.{}",
1623 remote_module_service_state_segment(&spec.module_name),
1624 remote_module_service_state_segment(&spec.service_name),
1625 extension
1626 ))
1627}
1628
1629fn remote_module_service_state_segment(value: &str) -> String {
1630 let mut segment = String::new();
1631 let mut previous_dash = false;
1632 for character in value.chars() {
1633 if character.is_ascii_alphanumeric() {
1634 segment.push(character.to_ascii_lowercase());
1635 previous_dash = false;
1636 } else if !segment.is_empty() && !previous_dash {
1637 segment.push('-');
1638 previous_dash = true;
1639 }
1640 }
1641 while segment.ends_with('-') {
1642 segment.pop();
1643 }
1644 if segment.is_empty() {
1645 "service".to_owned()
1646 } else {
1647 segment
1648 }
1649}
1650
1651fn terminate_remote_module_service(child: &mut Child) {
1652 if matches!(child.try_wait(), Ok(Some(_))) {
1653 return;
1654 }
1655
1656 #[cfg(unix)]
1657 {
1658 let process_group_id = child.id();
1659 let _ = Command::new("kill")
1660 .arg("-TERM")
1661 .arg(format!("-{process_group_id}"))
1662 .status();
1663 if wait_for_remote_module_service_exit(
1664 child,
1665 Duration::from_millis(REMOTE_SERVICE_TERMINATE_GRACE_MS),
1666 ) {
1667 return;
1668 }
1669 }
1670
1671 let _ = child.kill();
1672 let _ = child.wait();
1673}
1674
1675fn wait_for_remote_module_service_exit(child: &mut Child, timeout: Duration) -> bool {
1676 let started = Instant::now();
1677 loop {
1678 match child.try_wait() {
1679 Ok(Some(_)) => return true,
1680 Ok(None) => {}
1681 Err(_) => return true,
1682 }
1683 if started.elapsed() >= timeout {
1684 return false;
1685 }
1686 thread::sleep(Duration::from_millis(50));
1687 }
1688}
1689
1690#[cfg(unix)]
1691fn configure_remote_module_service_process(command: &mut Command) {
1692 use std::os::unix::process::CommandExt;
1693 command.process_group(0);
1694}
1695
1696#[cfg(not(unix))]
1697fn configure_remote_module_service_process(_command: &mut Command) {}
1698
1699fn shell_command(command: &str) -> Command {
1700 if cfg!(windows) {
1701 let mut process = Command::new("cmd");
1702 process.arg("/C").arg(command);
1703 process
1704 } else {
1705 let mut process = Command::new("sh");
1706 process.arg("-c").arg(command);
1707 process
1708 }
1709}
1710
1711fn current_timestamp() -> String {
1712 use platform_core::Clock;
1713 platform_core::SystemClock.now().to_rfc3339()
1714}
1715
1716fn duration_ms(started: Instant) -> u64 {
1717 u64::try_from(started.elapsed().as_millis()).unwrap_or(u64::MAX)
1718}
1719
1720#[must_use]
1722pub fn function_registry(modules: &[Module]) -> FunctionRegistry {
1723 let mut registry = FunctionRegistry::default();
1724 for module in modules {
1725 module.binding.register_functions(&mut registry);
1726 }
1727 registry
1728}
1729
1730pub async fn enqueue_lifecycle_activation_jobs(
1736 ctx: &AppContext,
1737 modules: &[Module],
1738 registry: &FunctionRegistry,
1739) -> platform_core::AppResult<Vec<String>> {
1740 validate_lifecycle_activation_jobs(modules, registry)?;
1741
1742 let client = RuntimeClient::new(ctx.db.clone());
1743 let mut run_ids = Vec::new();
1744
1745 for module in modules {
1746 let Some(lifecycle) = &module.manifest.lifecycle else {
1747 continue;
1748 };
1749
1750 for job in &lifecycle.activation_jobs {
1751 if job.run_policy != LifecycleActivationRunPolicy::EveryStartup {
1752 continue;
1753 }
1754 if !module_declares_runtime_function(module, &job.function_name) {
1755 continue;
1756 }
1757
1758 let Some(definition) = registry.get(&job.function_name) else {
1759 continue;
1760 };
1761
1762 let enqueue_result = client
1763 .enqueue_function(EnqueueFunctionRequest {
1764 function_name: job.function_name.clone(),
1765 input_json: job.input.clone(),
1766 correlation_id: CorrelationId::new(ctx.ids.new_id("corr_lifecycle")),
1767 actor: ActorContext::Service {
1768 service_id: "worker".to_owned(),
1769 scopes: vec!["runtime.functions.enqueue".to_owned()],
1770 },
1771 trace: TraceContext::default(),
1772 causation_id: Some(format!(
1773 "module_lifecycle:{}:{}",
1774 module.manifest.name, job.name
1775 )),
1776 max_attempts: Some(runtime_max_attempts_for_enqueue(
1777 definition.retry_policy.max_attempts,
1778 )),
1779 })
1780 .await;
1781
1782 match enqueue_result {
1783 Ok(run_id) => run_ids.push(run_id),
1784 Err(error) if job.required => return Err(error),
1785 Err(error) => warn_optional_lifecycle_enqueue_failure(
1786 &module.manifest.name,
1787 &job.name,
1788 &job.function_name,
1789 &error,
1790 ),
1791 }
1792 }
1793 }
1794
1795 Ok(run_ids)
1796}
1797
1798fn validate_lifecycle_activation_jobs(
1799 modules: &[Module],
1800 registry: &FunctionRegistry,
1801) -> platform_core::AppResult<()> {
1802 for module in modules {
1803 let Some(lifecycle) = &module.manifest.lifecycle else {
1804 continue;
1805 };
1806
1807 for check in &lifecycle.startup_checks {
1808 match &check.check {
1809 LifecycleStartupCheckKind::FunctionRegistered { function_name } => {
1810 if !module_declares_runtime_function(module, function_name) {
1811 let reason = format!(
1812 "startup check `{}` references function `{}` not declared by module `{}`",
1813 check.name, function_name, module.manifest.name
1814 );
1815 if !check.required {
1816 warn_optional_lifecycle_skip(
1817 &module.manifest.name,
1818 "startup_checks",
1819 &check.name,
1820 &reason,
1821 );
1822 continue;
1823 }
1824 return Err(lifecycle_validation_error(
1825 &module.manifest.name,
1826 "startup_checks",
1827 &check.name,
1828 format!("required {reason}"),
1829 ));
1830 }
1831 if registry.get(function_name).is_none() {
1832 let reason = format!(
1833 "startup check `{}` references missing function `{}`",
1834 check.name, function_name
1835 );
1836 if !check.required {
1837 warn_optional_lifecycle_skip(
1838 &module.manifest.name,
1839 "startup_checks",
1840 &check.name,
1841 &reason,
1842 );
1843 continue;
1844 }
1845 return Err(lifecycle_validation_error(
1846 &module.manifest.name,
1847 "startup_checks",
1848 &check.name,
1849 format!("required {reason}"),
1850 ));
1851 }
1852 }
1853 LifecycleStartupCheckKind::CapabilityDeclared { capability } => {
1854 if !module.manifest.capabilities.contains(capability) {
1855 let reason = format!(
1856 "startup check `{}` references missing capability `{}`",
1857 check.name, capability
1858 );
1859 if !check.required {
1860 warn_optional_lifecycle_skip(
1861 &module.manifest.name,
1862 "startup_checks",
1863 &check.name,
1864 &reason,
1865 );
1866 continue;
1867 }
1868 return Err(lifecycle_validation_error(
1869 &module.manifest.name,
1870 "startup_checks",
1871 &check.name,
1872 format!("required {reason}"),
1873 ));
1874 }
1875 }
1876 _ => {
1877 let reason = format!(
1878 "startup check `{}` uses an unsupported lifecycle check kind",
1879 check.name
1880 );
1881 if !check.required {
1882 warn_optional_lifecycle_skip(
1883 &module.manifest.name,
1884 "startup_checks",
1885 &check.name,
1886 &reason,
1887 );
1888 continue;
1889 }
1890 return Err(lifecycle_validation_error(
1891 &module.manifest.name,
1892 "startup_checks",
1893 &check.name,
1894 format!("required {reason}"),
1895 ));
1896 }
1897 }
1898 }
1899
1900 for job in &lifecycle.activation_jobs {
1901 if job.run_policy != LifecycleActivationRunPolicy::EveryStartup {
1902 continue;
1903 }
1904
1905 if !module_declares_runtime_function(module, &job.function_name) {
1906 let reason = format!(
1907 "activation job `{}` references function `{}` not declared by module `{}`",
1908 job.name, job.function_name, module.manifest.name
1909 );
1910 if !job.required {
1911 warn_optional_lifecycle_skip(
1912 &module.manifest.name,
1913 "activation_jobs",
1914 &job.name,
1915 &reason,
1916 );
1917 continue;
1918 }
1919 return Err(lifecycle_validation_error(
1920 &module.manifest.name,
1921 "activation_jobs",
1922 &job.name,
1923 format!("required {reason}"),
1924 ));
1925 }
1926 if registry.get(&job.function_name).is_none() {
1927 let reason = format!(
1928 "activation job `{}` references missing function `{}`",
1929 job.name, job.function_name
1930 );
1931 if !job.required {
1932 warn_optional_lifecycle_skip(
1933 &module.manifest.name,
1934 "activation_jobs",
1935 &job.name,
1936 &reason,
1937 );
1938 continue;
1939 }
1940 return Err(lifecycle_validation_error(
1941 &module.manifest.name,
1942 "activation_jobs",
1943 &job.name,
1944 format!("required {reason}"),
1945 ));
1946 }
1947 }
1948 }
1949
1950 Ok(())
1951}
1952
1953fn module_declares_runtime_function(module: &Module, function_name: &str) -> bool {
1954 module.manifest.runtime.as_ref().is_some_and(|runtime| {
1955 runtime
1956 .functions
1957 .iter()
1958 .any(|function| function.name == function_name)
1959 })
1960}
1961
1962fn lifecycle_validation_error(
1963 module_name: &str,
1964 collection: &str,
1965 item_name: &str,
1966 reason: String,
1967) -> AppError {
1968 AppError::validation(
1969 "Module lifecycle declaration failed validation",
1970 vec![ErrorDetail {
1971 field: Some(format!(
1972 "module.{module_name}.lifecycle.{collection}.{item_name}"
1973 )),
1974 reason,
1975 }],
1976 )
1977}
1978
1979fn warn_optional_lifecycle_skip(
1980 module_name: &str,
1981 collection: &str,
1982 item_name: &str,
1983 reason: &str,
1984) {
1985 tracing::warn!(
1986 module_name = %module_name,
1987 lifecycle_collection = %collection,
1988 lifecycle_item = %item_name,
1989 reason = %reason,
1990 "optional module lifecycle declaration skipped"
1991 );
1992}
1993
1994fn warn_optional_lifecycle_enqueue_failure(
1995 module_name: &str,
1996 job_name: &str,
1997 function_name: &str,
1998 error: &AppError,
1999) {
2000 tracing::warn!(
2001 module_name = %module_name,
2002 lifecycle_collection = "activation_jobs",
2003 lifecycle_item = %job_name,
2004 function_name = %function_name,
2005 error_code = %error.code.as_str(),
2006 error_message = %error.public_message,
2007 "optional module lifecycle activation enqueue failed"
2008 );
2009}
2010
2011fn runtime_max_attempts_for_enqueue(max_attempts: u32) -> i32 {
2012 i32::try_from(max_attempts).unwrap_or(i32::MAX)
2013}
2014
2015#[must_use]
2017pub fn event_handlers(modules: &[Module]) -> EventHandlerRegistry {
2018 event_handlers_with_context(modules, &EventHandlerRegistrationContext::empty())
2019}
2020
2021#[must_use]
2024pub fn event_handlers_with_runtime_actions(
2025 ctx: &AppContext,
2026 modules: &[Module],
2027 function_registry: Arc<FunctionRegistry>,
2028) -> EventHandlerRegistry {
2029 let context = EventHandlerRegistrationContext::with_runtime(
2030 RuntimeClient::new(ctx.db.clone()),
2031 function_registry,
2032 );
2033 event_handlers_with_context(modules, &context)
2034}
2035
2036fn event_handlers_with_context(
2037 modules: &[Module],
2038 context: &EventHandlerRegistrationContext,
2039) -> EventHandlerRegistry {
2040 let mut registry = EventHandlerRegistry::new();
2041 for module in modules {
2042 module
2043 .binding
2044 .register_event_handlers(&mut registry, context);
2045 }
2046 registry
2047}
2048
2049pub fn merge_linked_http(base: ApiOpenApiRouter) -> ApiOpenApiRouter {
2057 merge_linked_http_for_profile(base, CompositionProfile::default())
2058}
2059
2060pub fn merge_linked_http_for_profile(
2061 base: ApiOpenApiRouter,
2062 profile: CompositionProfile,
2063) -> ApiOpenApiRouter {
2064 linked_http_modules_for_profile(profile)
2065 .into_iter()
2066 .filter_map(|module| module.linked_http)
2067 .fold(base, |router, contribution| (contribution.merge)(router))
2068}
2069
2070pub fn merge_linked_http_for_config(
2071 base: ApiOpenApiRouter,
2072 config: &platform_core::AppConfig,
2073) -> platform_core::AppResult<ApiOpenApiRouter> {
2074 Ok(linked_http_modules_for_config(config)?
2075 .into_iter()
2076 .filter_map(|module| module.linked_http)
2077 .fold(base, |router, contribution| (contribution.merge)(router)))
2078}
2079
2080pub fn merge_linked_http_for_context(
2081 base: ApiOpenApiRouter,
2082 ctx: &AppContext,
2083) -> platform_core::AppResult<ApiOpenApiRouter> {
2084 Ok(linked_http_modules_for_context(ctx)?
2085 .into_iter()
2086 .filter_map(|module| module.linked_http)
2087 .fold(base, |router, contribution| (contribution.merge)(router)))
2088}
2089
2090pub fn merge_linked_http_for_context_with_composition(
2091 base: ApiOpenApiRouter,
2092 ctx: &AppContext,
2093 composition: &HostComposition,
2094) -> platform_core::AppResult<ApiOpenApiRouter> {
2095 Ok(
2096 linked_http_modules_for_context_with_composition(ctx, composition)?
2097 .into_iter()
2098 .filter_map(|module| module.linked_http)
2099 .fold(base, |router, contribution| (contribution.merge)(router)),
2100 )
2101}
2102
2103#[must_use]
2106pub fn story_display_descriptors() -> Vec<StoryDisplayDescriptor> {
2107 story_display_descriptors_for_profile(CompositionProfile::default())
2108}
2109
2110#[must_use]
2111pub fn story_display_descriptors_for_profile(
2112 profile: CompositionProfile,
2113) -> Vec<StoryDisplayDescriptor> {
2114 module_manifests_for_profile(profile)
2115 .into_iter()
2116 .flat_map(story_display_descriptors_from_manifest)
2117 .collect()
2118}
2119
2120pub fn story_display_descriptors_for_config(
2121 config: &platform_core::AppConfig,
2122) -> platform_core::AppResult<Vec<StoryDisplayDescriptor>> {
2123 Ok(linked_module_entries_for_config(config)?
2124 .into_iter()
2125 .flat_map(|entry| story_display_descriptors_from_manifest((entry.manifest)()))
2126 .collect())
2127}
2128
2129pub fn story_display_descriptors_for_context(
2130 ctx: &AppContext,
2131) -> platform_core::AppResult<Vec<StoryDisplayDescriptor>> {
2132 Ok(linked_module_entries_for_context(ctx)?
2133 .into_iter()
2134 .flat_map(|entry| story_display_descriptors_from_manifest((entry.manifest)()))
2135 .collect())
2136}
2137
2138pub fn install_default_story_display_catalog(ctx: &AppContext) -> platform_core::AppResult<()> {
2139 install_default_story_display_catalog_with_composition(ctx, &HostComposition::default())
2140}
2141
2142pub fn install_default_story_display_catalog_with_composition(
2143 ctx: &AppContext,
2144 composition: &HostComposition,
2145) -> platform_core::AppResult<()> {
2146 if !linked_module_enabled(ctx, story::module::MODULE_NAME) {
2147 story::backend::install_default_story_display(Vec::new());
2148 return Ok(());
2149 }
2150 let mut descriptors = story_display_descriptors_for_context(ctx)?;
2151 descriptors.extend(
2152 host_linked_modules_for_context(ctx, composition)
2153 .into_iter()
2154 .flat_map(|entry| story_display_descriptors_from_manifest((entry.manifest)())),
2155 );
2156 story::backend::install_default_story_display(descriptors);
2157 Ok(())
2158}
2159
2160pub fn install_story_display_catalog(metadata: &[AdminModuleMetadata]) {
2161 story::backend::install_story_display(
2162 metadata
2163 .iter()
2164 .flat_map(story_display_descriptors_from_metadata)
2165 .collect(),
2166 );
2167}
2168
2169fn story_display_descriptors_from_metadata(
2170 module: &AdminModuleMetadata,
2171) -> Vec<StoryDisplayDescriptor> {
2172 story_display_descriptors_from_manifest(
2173 ModuleManifest::builder(module.module_name.clone())
2174 .story_display(module.story_display.clone())
2175 .http_routes(module.http_routes.clone())
2176 .build(),
2177 )
2178}
2179
2180fn story_display_descriptors_from_manifest(
2181 manifest: ModuleManifest,
2182) -> Vec<StoryDisplayDescriptor> {
2183 let mut descriptors = manifest.story_display;
2184 let existing_http = descriptors
2185 .iter()
2186 .filter_map(|descriptor| match &descriptor.source {
2187 StoryDisplaySource::HttpRequest { method, path } => {
2188 Some((method.clone(), path.clone()))
2189 }
2190 StoryDisplaySource::ExecutionName { .. } => None,
2191 })
2192 .collect::<Vec<_>>();
2193
2194 descriptors.extend(manifest.http_routes.into_iter().filter_map(|route| {
2195 let display_name = route.display_name?;
2196 let method = http_method_label(route.method)?;
2197 if existing_http
2198 .iter()
2199 .any(|(existing_method, existing_path)| {
2200 existing_method == method && existing_path == &route.path
2201 })
2202 {
2203 return None;
2204 }
2205 Some(StoryDisplayDescriptor {
2206 source: StoryDisplaySource::HttpRequest {
2207 method: method.to_owned(),
2208 path: route.path,
2209 },
2210 display_name,
2211 story_title: route.story_title,
2212 })
2213 }));
2214 descriptors
2215}
2216
2217fn http_method_label(method: ModuleHttpMethod) -> Option<&'static str> {
2218 Some(match method {
2219 ModuleHttpMethod::Get => "GET",
2220 ModuleHttpMethod::Post => "POST",
2221 ModuleHttpMethod::Put => "PUT",
2222 ModuleHttpMethod::Patch => "PATCH",
2223 ModuleHttpMethod::Delete => "DELETE",
2224 _ => return None,
2225 })
2226}
2227
2228pub fn runtime_config_descriptors(
2233 ctx: &AppContext,
2234) -> platform_core::AppResult<Vec<RuntimeConfigDescriptor>> {
2235 runtime_config_descriptors_with_composition(ctx, &HostComposition::default())
2236}
2237
2238pub fn runtime_config_descriptors_with_composition(
2239 ctx: &AppContext,
2240 composition: &HostComposition,
2241) -> platform_core::AppResult<Vec<RuntimeConfigDescriptor>> {
2242 let profile = CompositionProfile::from_config(&ctx.config)?;
2243 let module_enabled_descriptors =
2244 linked_module_entries(profile)
2245 .iter()
2246 .map(|entry| RuntimeConfigDescriptor {
2247 key: module_enabled_config_key(entry.module_name),
2248 scope: RuntimeConfigScope::Shared,
2249 group: Some("modules"),
2250 section: None,
2251 order: 10,
2252 visible_when: None,
2253 generated: None,
2254 value_type: RuntimeConfigType::Bool,
2255 default: serde_json::json!(linked_module_enabled_from_config(
2256 &ctx.config,
2257 entry.module_name
2258 )),
2259 editable: true,
2260 restart_only: true,
2261 description: "Whether this linked module is loaded on service startup.",
2262 });
2263 let host_module_enabled_descriptors =
2264 composition
2265 .linked_modules()
2266 .iter()
2267 .map(|entry| RuntimeConfigDescriptor {
2268 key: module_enabled_config_key(entry.module_name),
2269 scope: RuntimeConfigScope::Shared,
2270 group: Some("modules"),
2271 section: None,
2272 order: 10,
2273 visible_when: None,
2274 generated: None,
2275 value_type: RuntimeConfigType::Bool,
2276 default: serde_json::json!(linked_module_enabled_from_config(
2277 &ctx.config,
2278 entry.module_name
2279 )),
2280 editable: true,
2281 restart_only: true,
2282 description: "Whether this host linked module is loaded on service startup.",
2283 });
2284 let remote_module_enabled_descriptors =
2285 ctx.config
2286 .module_sources
2287 .remote
2288 .iter()
2289 .map(|source| RuntimeConfigDescriptor {
2290 key: module_enabled_config_key(&source.name),
2291 scope: RuntimeConfigScope::Shared,
2292 group: Some("modules"),
2293 section: None,
2294 order: 10,
2295 visible_when: None,
2296 generated: None,
2297 value_type: RuntimeConfigType::Bool,
2298 default: serde_json::json!(remote_module_enabled_from_config(
2299 &ctx.config,
2300 &source.name
2301 )),
2302 editable: true,
2303 restart_only: true,
2304 description: "Whether this remote module is loaded on service startup.",
2305 });
2306 let module_descriptors = linked_module_entries(profile)
2307 .iter()
2308 .filter(|entry| linked_module_enabled_from_config(&ctx.config, entry.module_name))
2309 .map(|entry| (entry.load)(ctx))
2310 .chain(
2311 host_linked_modules_for_config(&ctx.config, composition)
2312 .into_iter()
2313 .map(|entry| load_host_linked_module(ctx, entry)),
2314 )
2315 .flat_map(|module| module.runtime_config.iter().cloned())
2316 .collect::<Vec<_>>();
2317 Ok(platform_core::worker_runtime_config::RUNTIME_CONFIG
2320 .iter()
2321 .cloned()
2322 .chain(module_enabled_descriptors)
2323 .chain(host_module_enabled_descriptors)
2324 .chain(remote_module_enabled_descriptors)
2325 .chain(module_descriptors)
2326 .collect())
2327}
2328
2329pub fn runtime_config_group_descriptors(
2331 ctx: &AppContext,
2332) -> platform_core::AppResult<Vec<RuntimeConfigGroupDescriptor>> {
2333 runtime_config_group_descriptors_with_composition(ctx, &HostComposition::default())
2334}
2335
2336pub fn runtime_config_group_descriptors_with_composition(
2337 ctx: &AppContext,
2338 composition: &HostComposition,
2339) -> platform_core::AppResult<Vec<RuntimeConfigGroupDescriptor>> {
2340 let profile = CompositionProfile::from_config(&ctx.config)?;
2341 let module_groups = linked_module_entries(profile)
2342 .iter()
2343 .filter(|entry| linked_module_enabled_from_config(&ctx.config, entry.module_name))
2344 .map(|entry| (entry.load)(ctx))
2345 .chain(
2346 host_linked_modules_for_config(&ctx.config, composition)
2347 .into_iter()
2348 .map(|entry| load_host_linked_module(ctx, entry)),
2349 )
2350 .flat_map(|module| module.runtime_config_groups.iter().cloned())
2351 .collect::<Vec<_>>();
2352
2353 Ok(std::iter::once(MODULES_CONFIG_GROUP.clone())
2354 .chain(
2355 platform_core::worker_runtime_config::RUNTIME_CONFIG_GROUPS
2356 .iter()
2357 .cloned(),
2358 )
2359 .chain(module_groups)
2360 .collect())
2361}
2362
2363#[cfg(test)]
2364mod tests {
2365 use super::*;
2366 use async_trait::async_trait;
2367 use platform_core::{
2368 AppConfig, AuthConfig, DatabaseConfig, ErrorCode, ExecutionContext, HttpConfig,
2369 LoggingEventPublisher, ModuleConfig, ModuleSourcesConfig, PLATFORM_MIGRATIONS, RedisConfig,
2370 RemoteModuleSourceConfig, RuntimeConfigProvider, RuntimeConfigRegistry,
2371 RuntimeConfigSnapshot, ServiceConfig, TelemetryConfig, apply_migrations,
2372 };
2373 use platform_module::{
2374 ConsoleArea, LifecycleActivationJobDeclaration, LifecycleStartupCheckDeclaration,
2375 LifecycleSurface, ModuleManifestLintSeverity, RuntimeFunctionDeclaration, RuntimeSurface,
2376 lint_module_manifest,
2377 };
2378 use platform_runtime::{FunctionDefinition, FunctionHandler, RUNTIME_MIGRATIONS, RetryPolicy};
2379 use platform_testing::{SequentialIdGenerator, TestDatabase};
2380 use serde_json::{Value, json};
2381 use sqlx::postgres::{PgConnectOptions, PgPoolOptions};
2382 use std::collections::BTreeMap;
2383 use std::sync::Arc;
2384 use std::time::Duration;
2385
2386 #[derive(Debug)]
2387 struct TestRuntimeConfigProvider {
2388 snapshot: Arc<RuntimeConfigSnapshot>,
2389 }
2390
2391 impl RuntimeConfigProvider for TestRuntimeConfigProvider {
2392 fn snapshot(&self) -> Arc<RuntimeConfigSnapshot> {
2393 Arc::clone(&self.snapshot)
2394 }
2395 }
2396
2397 #[test]
2398 fn linked_module_entry_names_match_manifests() {
2399 for profile in [CompositionProfile::Core, CompositionProfile::Demo] {
2400 for entry in linked_module_entries(profile) {
2401 assert_eq!(
2402 entry.module_name,
2403 (entry.manifest)().name,
2404 "linked module entry name must match ModuleManifest::name"
2405 );
2406 }
2407 }
2408 }
2409
2410 #[test]
2411 fn core_profile_excludes_demo_linked_modules() {
2412 let names = module_manifests_for_profile(CompositionProfile::Core)
2413 .into_iter()
2414 .map(|manifest| manifest.name)
2415 .collect::<Vec<_>>();
2416
2417 assert_eq!(names, vec!["platform-story"]);
2418 }
2419
2420 #[test]
2421 fn demo_profile_includes_fixture_linked_modules() {
2422 let names = module_manifests_for_profile(CompositionProfile::Demo)
2423 .into_iter()
2424 .map(|manifest| manifest.name)
2425 .collect::<Vec<_>>();
2426
2427 assert_eq!(names, vec!["auth", "auth-password", "platform-story"]);
2428 }
2429
2430 #[test]
2431 fn http_route_metadata_contributes_story_display_descriptors() {
2432 let descriptors = story_display_descriptors_for_profile(CompositionProfile::Demo);
2433
2434 assert!(descriptors.iter().any(|descriptor| {
2435 matches!(
2436 &descriptor.source,
2437 StoryDisplaySource::HttpRequest { method, path }
2438 if method == "POST" && path == "/v1/auth/dev/sessions"
2439 ) && descriptor.display_name == "Create Development Session"
2440 }));
2441 }
2442
2443 #[test]
2444 fn core_profile_migrations_exclude_demo_module_migrations() {
2445 let names = migrations_for_profile(CompositionProfile::Core)
2446 .into_iter()
2447 .map(|migration| migration.name)
2448 .collect::<Vec<_>>();
2449
2450 assert!(names.iter().any(|name| name.starts_with("platform/")));
2451 assert!(names.iter().any(|name| name.starts_with("runtime/")));
2452 assert!(!names.iter().any(|name| name.starts_with("auth/")));
2453 assert!(!names.iter().any(|name| name.starts_with("auth-password/")));
2454 }
2455
2456 #[test]
2457 fn demo_profile_migrations_include_fixture_module_migrations() {
2458 let names = migrations_for_profile(CompositionProfile::Demo)
2459 .into_iter()
2460 .map(|migration| migration.name)
2461 .collect::<Vec<_>>();
2462
2463 assert!(
2464 names
2465 .iter()
2466 .any(|name| name == &"auth/0001_create_auth_schema")
2467 );
2468 assert!(
2469 names
2470 .iter()
2471 .any(|name| name == &"auth-password/0001_create_auth_password_schema")
2472 );
2473 }
2474
2475 #[test]
2476 fn host_composition_migrations_include_enabled_host_linked_modules() {
2477 let config = test_config_with_database_url("postgres://localhost/lenso_test");
2478 let composition = HostComposition::new().with_linked_module(test_host_linked_module());
2479
2480 let names = migrations_for_config_with_composition(&config, &composition)
2481 .expect("host composition migrations should load")
2482 .into_iter()
2483 .map(|migration| migration.name)
2484 .collect::<Vec<_>>();
2485
2486 assert!(names.iter().any(|name| name == &"billing/0001_init"));
2487 }
2488
2489 #[test]
2490 fn host_composition_can_install_auth_modules() {
2491 let mut config = test_config_with_database_url("postgres://localhost/lenso_test");
2492 config.module_sources.linked_profile = "core".to_owned();
2493 let composition = HostComposition::new()
2494 .with_linked_module(auth_linked_module())
2495 .with_linked_module(auth_password_linked_module());
2496
2497 let names = migrations_for_config_with_composition(&config, &composition)
2498 .expect("host composition migrations should load")
2499 .into_iter()
2500 .map(|migration| migration.name)
2501 .collect::<Vec<_>>();
2502
2503 assert!(
2504 names
2505 .iter()
2506 .any(|name| name == &"auth/0001_create_auth_schema")
2507 );
2508 assert!(
2509 names
2510 .iter()
2511 .any(|name| name == &"auth-password/0001_create_auth_password_schema")
2512 );
2513 }
2514
2515 #[tokio::test]
2516 async fn host_composition_runtime_config_includes_host_module_toggle() {
2517 let db = platform_core::DbPool::connect_lazy("postgres://localhost/lenso_test")
2518 .expect("lazy pool should build");
2519 let config = test_config_with_database_url("postgres://localhost/lenso_test");
2520 let ctx = AppContext::new(config, db, Arc::new(LoggingEventPublisher));
2521 let composition = HostComposition::new().with_linked_module(test_host_linked_module());
2522
2523 let keys = runtime_config_descriptors_with_composition(&ctx, &composition)
2524 .expect("host composition descriptors should load")
2525 .into_iter()
2526 .map(|descriptor| descriptor.key)
2527 .collect::<Vec<_>>();
2528
2529 assert!(keys.iter().any(|key| key == "modules.billing.enabled"));
2530 }
2531
2532 #[tokio::test]
2533 async fn host_composition_modules_include_manifest_only_modules() {
2534 let db = platform_core::DbPool::connect_lazy("postgres://localhost/lenso_test")
2535 .expect("lazy pool should build");
2536 let config = test_config_with_database_url("postgres://localhost/lenso_test");
2537 let ctx = AppContext::new(config, db, Arc::new(LoggingEventPublisher));
2538 let composition = HostComposition::new().with_linked_module(test_host_linked_module());
2539
2540 let names = modules_for_config_with_composition(&ctx, &composition)
2541 .expect("host composition modules should load")
2542 .into_iter()
2543 .map(|module| module.manifest.name)
2544 .collect::<Vec<_>>();
2545
2546 assert!(names.iter().any(|name| name == "billing"));
2547 }
2548
2549 #[test]
2550 fn demo_profile_includes_every_core_entry() {
2551 let demo_names = linked_module_entries(CompositionProfile::Demo)
2552 .iter()
2553 .map(|entry| entry.module_name)
2554 .collect::<Vec<_>>();
2555
2556 for core_entry in linked_module_entries(CompositionProfile::Core) {
2557 assert!(
2558 demo_names.contains(&core_entry.module_name),
2559 "demo profile should include core linked module `{}`",
2560 core_entry.module_name
2561 );
2562 }
2563 }
2564
2565 #[test]
2566 fn default_module_manifests_use_demo_profile() {
2567 let names = module_manifests()
2568 .into_iter()
2569 .map(|manifest| manifest.name)
2570 .collect::<Vec<_>>();
2571
2572 assert_eq!(names, vec!["auth", "auth-password", "platform-story"]);
2573 }
2574
2575 #[test]
2576 fn linked_http_route_owners_are_profile_aware() {
2577 assert_eq!(
2578 linked_http_route_owners_for_profile(CompositionProfile::Core),
2579 vec![LinkedHttpRouteOwner {
2580 module_name: "platform-story".to_owned(),
2581 public_prefixes: &["/admin/runtime/stories"],
2582 }]
2583 );
2584 assert_eq!(
2585 linked_http_route_owners_for_profile(CompositionProfile::Demo),
2586 vec![
2587 LinkedHttpRouteOwner {
2588 module_name: "auth".to_owned(),
2589 public_prefixes: &["/v1/auth/dev/", "/v1/auth/sessions/"],
2590 },
2591 LinkedHttpRouteOwner {
2592 module_name: "auth-password".to_owned(),
2593 public_prefixes: &["/v1/auth/password/"],
2594 },
2595 LinkedHttpRouteOwner {
2596 module_name: "platform-story".to_owned(),
2597 public_prefixes: &["/admin/runtime/stories"],
2598 },
2599 ]
2600 );
2601 }
2602
2603 #[tokio::test]
2604 async fn modules_for_config_uses_core_linked_profile() {
2605 let db = platform_core::DbPool::connect_lazy("postgres://localhost/lenso_test")
2606 .expect("lazy pool should build");
2607 let mut config = test_config_with_database_url("postgres://localhost/lenso_test");
2608 config.module_sources.linked_profile = "core".to_owned();
2609 let ctx = AppContext::new(config, db, Arc::new(LoggingEventPublisher));
2610
2611 let names = modules_for_config(&ctx)
2612 .expect("core linked profile should parse")
2613 .into_iter()
2614 .map(|module| module.manifest.name)
2615 .collect::<Vec<_>>();
2616
2617 assert_eq!(names, vec!["platform-story"]);
2618 }
2619
2620 #[tokio::test]
2621 async fn auth_actor_resolver_is_profile_and_composition_aware() {
2622 let db = platform_core::DbPool::connect_lazy("postgres://localhost/lenso_test")
2623 .expect("lazy pool should build");
2624 let demo_ctx = AppContext::new(
2625 test_config_with_database_url("postgres://localhost/lenso_test"),
2626 db.clone(),
2627 Arc::new(LoggingEventPublisher),
2628 );
2629 assert!(
2630 auth_actor_resolver_for_context(&demo_ctx)
2631 .expect("demo profile")
2632 .is_some()
2633 );
2634
2635 let mut composition_config =
2636 test_config_with_database_url("postgres://localhost/lenso_test");
2637 composition_config.module_sources.linked_profile = "core".to_owned();
2638 let composition_ctx = AppContext::new(
2639 composition_config,
2640 db.clone(),
2641 Arc::new(LoggingEventPublisher),
2642 );
2643 let composition = HostComposition::new().with_linked_module(auth_linked_module());
2644 assert!(
2645 auth_actor_resolver_for_context_with_composition(&composition_ctx, &composition)
2646 .expect("auth composition")
2647 .is_some()
2648 );
2649
2650 let mut core_config = test_config_with_database_url("postgres://localhost/lenso_test");
2651 core_config.module_sources.linked_profile = "core".to_owned();
2652 let core_ctx = AppContext::new(core_config, db, Arc::new(LoggingEventPublisher));
2653 assert!(
2654 auth_actor_resolver_for_context(&core_ctx)
2655 .expect("core profile")
2656 .is_none()
2657 );
2658 }
2659
2660 #[tokio::test]
2661 async fn auth_actor_resolver_respects_disabled_auth_module() {
2662 let db = platform_core::DbPool::connect_lazy("postgres://localhost/lenso_test")
2663 .expect("lazy pool should build");
2664 let mut config = test_config_with_database_url("postgres://localhost/lenso_test");
2665 config.modules.insert(
2666 auth::module::MODULE_NAME.to_owned(),
2667 ModuleConfig {
2668 enabled: Some(false),
2669 values: BTreeMap::new(),
2670 },
2671 );
2672 let ctx = AppContext::new(config, db, Arc::new(LoggingEventPublisher));
2673
2674 assert!(
2675 auth_actor_resolver_for_context(&ctx)
2676 .expect("demo profile")
2677 .is_none()
2678 );
2679 }
2680
2681 #[tokio::test]
2682 async fn auth_password_requires_auth_module() {
2683 let db = platform_core::DbPool::connect_lazy("postgres://localhost/lenso_test")
2684 .expect("lazy pool should build");
2685 let mut config = test_config_with_database_url("postgres://localhost/lenso_test");
2686 config.modules.insert(
2687 auth::module::MODULE_NAME.to_owned(),
2688 ModuleConfig {
2689 enabled: Some(false),
2690 values: BTreeMap::new(),
2691 },
2692 );
2693 let ctx = AppContext::new(config, db, Arc::new(LoggingEventPublisher));
2694
2695 let names = modules_for_config(&ctx)
2696 .expect("demo profile")
2697 .into_iter()
2698 .map(|module| module.manifest.name)
2699 .collect::<Vec<_>>();
2700
2701 assert!(!names.iter().any(|name| name == "auth-password"));
2702 }
2703
2704 #[tokio::test]
2705 async fn auth_password_dependency_status_is_visible_in_metadata() {
2706 let db = platform_core::DbPool::connect_lazy("postgres://localhost/lenso_test")
2707 .expect("lazy pool should build");
2708 let mut config = test_config_with_database_url("postgres://localhost/lenso_test");
2709 config.modules.insert(
2710 auth::module::MODULE_NAME.to_owned(),
2711 ModuleConfig {
2712 enabled: Some(false),
2713 values: BTreeMap::new(),
2714 },
2715 );
2716 let ctx = AppContext::new(config, db, Arc::new(LoggingEventPublisher));
2717
2718 let metadata = load_admin_module_metadata(&ctx)
2719 .await
2720 .expect("module metadata should load");
2721 let auth_password = metadata
2722 .iter()
2723 .find(|module| module.module_name == "auth-password")
2724 .expect("dependency-disabled provider should remain visible in metadata");
2725
2726 assert_eq!(
2727 auth_password.dependencies,
2728 vec![auth::module::MODULE_NAME.to_owned()]
2729 );
2730 assert!(matches!(
2731 &auth_password.load_status,
2732 ModuleLoadStatus::Error { message }
2733 if message == "module dependency disabled: auth"
2734 ));
2735 }
2736
2737 #[tokio::test]
2738 async fn auth_actor_resolver_allows_jwt_strategy_without_secret() {
2739 let db = platform_core::DbPool::connect_lazy("postgres://localhost/lenso_test")
2740 .expect("lazy pool should build");
2741 let config = test_config_with_database_url("postgres://localhost/lenso_test");
2742 let ctx = AppContext::new(config, db, Arc::new(LoggingEventPublisher));
2743 let registry =
2744 RuntimeConfigRegistry::try_new(runtime_config_descriptors(&ctx).expect("descriptors"))
2745 .expect("registry");
2746 let mut stored = BTreeMap::new();
2747 stored.insert(
2748 ("*".to_owned(), "auth-password.token_strategy".to_owned()),
2749 json!("jwt"),
2750 );
2751 let snapshot = RuntimeConfigSnapshot::resolve(®istry, "api", &stored);
2752 let ctx = ctx.with_runtime_config_provider(Arc::new(TestRuntimeConfigProvider {
2753 snapshot: Arc::new(snapshot),
2754 }));
2755
2756 assert!(
2757 auth_actor_resolver_for_context(&ctx)
2758 .expect("JWT resolver should be skipped until jwt_secret is configured")
2759 .is_some()
2760 );
2761 }
2762
2763 #[tokio::test]
2764 async fn auth_actor_resolver_requires_redis_when_session_cache_is_redis() {
2765 let db = platform_core::DbPool::connect_lazy("postgres://localhost/lenso_test")
2766 .expect("lazy pool should build");
2767 let config = test_config_with_database_url("postgres://localhost/lenso_test");
2768 let ctx = AppContext::new(config, db, Arc::new(LoggingEventPublisher));
2769 let registry =
2770 RuntimeConfigRegistry::try_new(runtime_config_descriptors(&ctx).expect("descriptors"))
2771 .expect("registry");
2772 let mut stored = BTreeMap::new();
2773 stored.insert(
2774 ("*".to_owned(), "auth.session_cache".to_owned()),
2775 json!("redis"),
2776 );
2777 let snapshot = RuntimeConfigSnapshot::resolve(®istry, "api", &stored);
2778 let ctx = ctx.with_runtime_config_provider(Arc::new(TestRuntimeConfigProvider {
2779 snapshot: Arc::new(snapshot),
2780 }));
2781
2782 let error =
2783 auth_actor_resolver_for_context(&ctx).expect_err("redis cache should require Redis");
2784
2785 assert_eq!(error.code, ErrorCode::Validation);
2786 }
2787
2788 #[tokio::test]
2789 async fn modules_for_config_skips_disabled_linked_modules() {
2790 let db = platform_core::DbPool::connect_lazy("postgres://localhost/lenso_test")
2791 .expect("lazy pool should build");
2792 let mut config = test_config_with_database_url("postgres://localhost/lenso_test");
2793 config.modules.insert(
2794 "auth-password".to_owned(),
2795 ModuleConfig {
2796 enabled: Some(false),
2797 values: BTreeMap::new(),
2798 },
2799 );
2800 let ctx = AppContext::new(config, db, Arc::new(LoggingEventPublisher));
2801
2802 let names = modules_for_config(&ctx)
2803 .expect("demo linked profile should parse")
2804 .into_iter()
2805 .map(|module| module.manifest.name)
2806 .collect::<Vec<_>>();
2807
2808 assert_eq!(names, vec!["auth", "platform-story"]);
2809 }
2810
2811 #[tokio::test]
2812 async fn modules_for_config_uses_runtime_config_enabled_flag() {
2813 let db = platform_core::DbPool::connect_lazy("postgres://localhost/lenso_test")
2814 .expect("lazy pool should build");
2815 let config = test_config_with_database_url("postgres://localhost/lenso_test");
2816 let ctx = AppContext::new(config, db, Arc::new(LoggingEventPublisher));
2817 let registry =
2818 RuntimeConfigRegistry::try_new(runtime_config_descriptors(&ctx).expect("descriptors"))
2819 .expect("registry");
2820 let mut stored = BTreeMap::new();
2821 stored.insert(
2822 ("*".to_owned(), "modules.auth-password.enabled".to_owned()),
2823 json!(false),
2824 );
2825 let snapshot = RuntimeConfigSnapshot::resolve(®istry, "api", &stored);
2826 let ctx = ctx.with_runtime_config_provider(Arc::new(TestRuntimeConfigProvider {
2827 snapshot: Arc::new(snapshot),
2828 }));
2829
2830 let names = modules_for_config(&ctx)
2831 .expect("demo linked profile should parse")
2832 .into_iter()
2833 .map(|module| module.manifest.name)
2834 .collect::<Vec<_>>();
2835
2836 assert_eq!(names, vec!["auth", "platform-story"]);
2837 let linked_http_names = linked_http_modules_for_context(&ctx)
2838 .expect("linked HTTP modules should load")
2839 .into_iter()
2840 .map(|module| module.manifest.name)
2841 .collect::<Vec<_>>();
2842
2843 assert_eq!(linked_http_names, vec!["auth", "platform-story"]);
2844 }
2845
2846 #[tokio::test]
2847 async fn story_module_runtime_config_disables_backend_metadata() {
2848 let db = platform_core::DbPool::connect_lazy("postgres://localhost/lenso_test")
2849 .expect("lazy pool should build");
2850 let config = test_config_with_database_url("postgres://localhost/lenso_test");
2851 let ctx = AppContext::new(config, db, Arc::new(LoggingEventPublisher));
2852 let registry =
2853 RuntimeConfigRegistry::try_new(runtime_config_descriptors(&ctx).expect("descriptors"))
2854 .expect("registry");
2855 let mut stored = BTreeMap::new();
2856 stored.insert(
2857 ("*".to_owned(), "modules.platform-story.enabled".to_owned()),
2858 json!(false),
2859 );
2860 let snapshot = RuntimeConfigSnapshot::resolve(®istry, "api", &stored);
2861 let ctx = ctx.with_runtime_config_provider(Arc::new(TestRuntimeConfigProvider {
2862 snapshot: Arc::new(snapshot),
2863 }));
2864
2865 let linked_http_names = linked_http_modules_for_context(&ctx)
2866 .expect("linked HTTP modules should load")
2867 .into_iter()
2868 .map(|module| module.manifest.name)
2869 .collect::<Vec<_>>();
2870 assert_eq!(linked_http_names, vec!["auth", "auth-password"]);
2871
2872 let metadata = load_admin_module_metadata(&ctx)
2873 .await
2874 .expect("module metadata should load");
2875 let story = metadata
2876 .iter()
2877 .find(|module| module.module_name == "platform-story")
2878 .expect("disabled story module should remain visible in metadata");
2879
2880 assert!(matches!(
2881 &story.load_status,
2882 ModuleLoadStatus::Error { message }
2883 if message == "module disabled by configuration"
2884 ));
2885 assert_eq!(story.console.len(), 1);
2886 assert_eq!(story.http_routes.len(), story::module::http_routes().len());
2887 }
2888
2889 #[tokio::test]
2890 async fn runtime_config_descriptors_include_module_enabled_flags() {
2891 let db = platform_core::DbPool::connect_lazy("postgres://localhost/lenso_test")
2892 .expect("lazy pool should build");
2893 let config = test_config_with_database_url("postgres://localhost/lenso_test");
2894 let ctx = AppContext::new(config, db, Arc::new(LoggingEventPublisher));
2895
2896 let keys = runtime_config_descriptors(&ctx)
2897 .expect("descriptors should load")
2898 .into_iter()
2899 .map(|descriptor| {
2900 (
2901 descriptor.key,
2902 descriptor.group,
2903 descriptor.restart_only,
2904 descriptor.default,
2905 )
2906 })
2907 .collect::<Vec<_>>();
2908
2909 assert!(keys.iter().any(|(key, group, restart_only, default)| {
2910 key == "modules.auth.enabled"
2911 && *group == Some("modules")
2912 && *restart_only
2913 && default == &json!(true)
2914 }));
2915 assert!(keys.iter().any(|(key, group, restart_only, default)| {
2916 key == "modules.auth-password.enabled"
2917 && *group == Some("modules")
2918 && *restart_only
2919 && default == &json!(true)
2920 }));
2921 assert!(keys.iter().any(|(key, group, restart_only, default)| {
2922 key == "modules.platform-story.enabled"
2923 && *group == Some("modules")
2924 && *restart_only
2925 && default == &json!(true)
2926 }));
2927 }
2928
2929 #[tokio::test]
2930 async fn runtime_config_groups_include_module_owned_groups() {
2931 let db = platform_core::DbPool::connect_lazy("postgres://localhost/lenso_test")
2932 .expect("lazy pool should build");
2933 let config = test_config_with_database_url("postgres://localhost/lenso_test");
2934 let ctx = AppContext::new(config, db, Arc::new(LoggingEventPublisher));
2935
2936 let groups = runtime_config_group_descriptors(&ctx)
2937 .expect("groups should load")
2938 .into_iter()
2939 .map(|group| (group.id, group.label))
2940 .collect::<Vec<_>>();
2941
2942 assert!(groups.contains(&("modules", "Modules")));
2943 assert!(groups.contains(&("auth-password.hashing", "Password Hashing")));
2944 assert!(groups.contains(&("auth-password.tokens", "Tokens")));
2945 assert!(!groups.iter().any(|(id, _)| *id == "auth-password.jwt"));
2946 }
2947
2948 #[tokio::test]
2949 async fn runtime_config_descriptors_include_remote_module_enabled_flags() {
2950 let db = platform_core::DbPool::connect_lazy("postgres://localhost/lenso_test")
2951 .expect("lazy pool should build");
2952 let mut config = test_config_with_database_url("postgres://localhost/lenso_test");
2953 config.module_sources.remote.push(RemoteModuleSourceConfig {
2954 name: "remote-crm".to_owned(),
2955 base_url: "http://127.0.0.1:65535".to_owned(),
2956 auth_token_env: None,
2957 timeout_ms: 1,
2958 });
2959 config.modules.insert(
2960 "remote-crm".to_owned(),
2961 ModuleConfig {
2962 enabled: Some(false),
2963 values: BTreeMap::new(),
2964 },
2965 );
2966 let ctx = AppContext::new(config, db, Arc::new(LoggingEventPublisher));
2967
2968 let keys = runtime_config_descriptors(&ctx)
2969 .expect("descriptors should load")
2970 .into_iter()
2971 .map(|descriptor| (descriptor.key, descriptor.restart_only, descriptor.default))
2972 .collect::<Vec<_>>();
2973
2974 assert!(keys.iter().any(|(key, restart_only, default)| {
2975 key == "modules.remote-crm.enabled" && *restart_only && default == &json!(false)
2976 }));
2977 }
2978
2979 #[tokio::test]
2980 async fn load_modules_skips_runtime_disabled_remote_modules() {
2981 let db = platform_core::DbPool::connect_lazy("postgres://localhost/lenso_test")
2982 .expect("lazy pool should build");
2983 let mut config = test_config_with_database_url("postgres://localhost/lenso_test");
2984 config.module_sources.remote.push(RemoteModuleSourceConfig {
2985 name: "remote-crm".to_owned(),
2986 base_url: "http://127.0.0.1:65535".to_owned(),
2987 auth_token_env: None,
2988 timeout_ms: 1,
2989 });
2990 let ctx = AppContext::new(config, db, Arc::new(LoggingEventPublisher));
2991 let registry =
2992 RuntimeConfigRegistry::try_new(runtime_config_descriptors(&ctx).expect("descriptors"))
2993 .expect("registry");
2994 let mut stored = BTreeMap::new();
2995 stored.insert(
2996 ("*".to_owned(), "modules.remote-crm.enabled".to_owned()),
2997 json!(false),
2998 );
2999 let snapshot = RuntimeConfigSnapshot::resolve(®istry, "api", &stored);
3000 let ctx = ctx.with_runtime_config_provider(Arc::new(TestRuntimeConfigProvider {
3001 snapshot: Arc::new(snapshot),
3002 }));
3003
3004 let names = load_modules(&ctx)
3005 .await
3006 .expect("disabled remote should not be loaded")
3007 .into_iter()
3008 .map(|module| module.manifest.name)
3009 .collect::<Vec<_>>();
3010
3011 assert!(!names.iter().any(|name| name == "remote-crm"));
3012 }
3013
3014 #[tokio::test]
3015 async fn module_metadata_reports_disabled_remote_modules() {
3016 let db = platform_core::DbPool::connect_lazy("postgres://localhost/lenso_test")
3017 .expect("lazy pool should build");
3018 let mut config = test_config_with_database_url("postgres://localhost/lenso_test");
3019 config.module_sources.remote.push(RemoteModuleSourceConfig {
3020 name: "remote-grpc-crm".to_owned(),
3021 base_url: "grpc://127.0.0.1:65535".to_owned(),
3022 auth_token_env: None,
3023 timeout_ms: 1,
3024 });
3025 config.modules.insert(
3026 "remote-grpc-crm".to_owned(),
3027 ModuleConfig {
3028 enabled: Some(false),
3029 values: BTreeMap::new(),
3030 },
3031 );
3032 let ctx = AppContext::new(config, db, Arc::new(LoggingEventPublisher));
3033
3034 let metadata = load_admin_module_metadata(&ctx)
3035 .await
3036 .expect("module metadata should load");
3037 let remote = metadata
3038 .iter()
3039 .find(|module| module.module_name == "remote-grpc-crm")
3040 .expect("disabled remote module should remain visible in metadata");
3041
3042 assert_eq!(remote.source, ModuleSource::Remote);
3043 assert!(matches!(
3044 &remote.load_status,
3045 ModuleLoadStatus::Error { message }
3046 if message == "module disabled by configuration"
3047 ));
3048 assert!(matches!(
3049 &remote.source_diagnostics,
3050 Some(AdminModuleSourceDiagnostics::Remote(diagnostics))
3051 if diagnostics.transport == "grpc"
3052 && diagnostics.base_url == "http://127.0.0.1:65535"
3053 ));
3054 }
3055
3056 #[test]
3057 fn migrations_for_config_skip_disabled_linked_module_migrations() {
3058 let mut config = test_config_with_database_url("postgres://localhost/lenso_test");
3059 config.modules.insert(
3060 "auth-password".to_owned(),
3061 ModuleConfig {
3062 enabled: Some(false),
3063 values: BTreeMap::new(),
3064 },
3065 );
3066
3067 let names = migrations_for_config(&config)
3068 .expect("demo linked profile should parse")
3069 .into_iter()
3070 .map(|migration| migration.name)
3071 .collect::<Vec<_>>();
3072
3073 assert!(!names.iter().any(|name| name.starts_with("auth-password/")));
3074 assert!(
3075 names
3076 .iter()
3077 .any(|name| name == &"auth/0001_create_auth_schema")
3078 );
3079 }
3080
3081 #[test]
3082 fn linked_http_modules_for_config_skip_disabled_linked_routes() {
3083 let mut config = test_config_with_database_url("postgres://localhost/lenso_test");
3084 config.modules.insert(
3085 "auth-password".to_owned(),
3086 ModuleConfig {
3087 enabled: Some(false),
3088 values: BTreeMap::new(),
3089 },
3090 );
3091
3092 let names = linked_http_modules_for_config(&config)
3093 .expect("demo linked profile should parse")
3094 .into_iter()
3095 .map(|module| module.manifest.name)
3096 .collect::<Vec<_>>();
3097
3098 assert_eq!(names, vec!["auth", "platform-story"]);
3099 }
3100
3101 #[test]
3102 fn linked_http_modules_for_config_skip_disabled_story_routes() {
3103 let mut config = test_config_with_database_url("postgres://localhost/lenso_test");
3104 config.modules.insert(
3105 "platform-story".to_owned(),
3106 ModuleConfig {
3107 enabled: Some(false),
3108 values: BTreeMap::new(),
3109 },
3110 );
3111
3112 let names = linked_http_modules_for_config(&config)
3113 .expect("demo linked profile should parse")
3114 .into_iter()
3115 .map(|module| module.manifest.name)
3116 .collect::<Vec<_>>();
3117
3118 assert_eq!(names, vec!["auth", "auth-password"]);
3119 }
3120
3121 #[tokio::test]
3122 async fn disabled_story_module_omits_default_story_display_catalog() {
3123 story::backend::reset_catalogs_for_test();
3124 let mut config = test_config_with_database_url("postgres://localhost/lenso_test");
3125 config.modules.insert(
3126 "platform-story".to_owned(),
3127 ModuleConfig {
3128 enabled: Some(false),
3129 values: BTreeMap::new(),
3130 },
3131 );
3132 let db = platform_core::DbPool::connect_lazy("postgres://localhost/lenso_test")
3133 .expect("lazy pool should build");
3134 let ctx = AppContext::new(config, db, Arc::new(LoggingEventPublisher));
3135
3136 install_default_story_display_catalog(&ctx)
3137 .expect("story display catalog installation should succeed");
3138
3139 assert!(story::backend::story_display_catalog_snapshot().is_empty());
3140 }
3141
3142 #[tokio::test]
3143 async fn module_metadata_reports_disabled_linked_modules() {
3144 let db = platform_core::DbPool::connect_lazy("postgres://localhost/lenso_test")
3145 .expect("lazy pool should build");
3146 let mut config = test_config_with_database_url("postgres://localhost/lenso_test");
3147 config.modules.insert(
3148 "auth-password".to_owned(),
3149 ModuleConfig {
3150 enabled: Some(false),
3151 values: BTreeMap::new(),
3152 },
3153 );
3154 let ctx = AppContext::new(config, db, Arc::new(LoggingEventPublisher));
3155
3156 let metadata = load_admin_module_metadata(&ctx)
3157 .await
3158 .expect("module metadata should load");
3159 let auth_password = metadata
3160 .iter()
3161 .find(|module| module.module_name == "auth-password")
3162 .expect("disabled module should remain visible in metadata");
3163
3164 assert!(matches!(
3165 &auth_password.load_status,
3166 ModuleLoadStatus::Error { message }
3167 if message == "module disabled by configuration"
3168 ));
3169 }
3170
3171 #[test]
3172 fn composition_profile_rejects_unknown_values() {
3173 let error = CompositionProfile::parse("fixture")
3174 .expect_err("fixture is not a supported linked module profile");
3175
3176 assert_eq!(error.code, ErrorCode::Validation);
3177 assert!(
3178 error
3179 .details
3180 .iter()
3181 .any(|detail| detail.field.as_deref() == Some("module_sources.linked_profile"))
3182 );
3183 }
3184
3185 #[test]
3186 fn linked_http_route_owners_are_projected_from_modules() {
3187 assert_eq!(
3188 linked_http_route_owners(),
3189 vec![
3190 LinkedHttpRouteOwner {
3191 module_name: "auth".to_owned(),
3192 public_prefixes: &["/v1/auth/dev/", "/v1/auth/sessions/"],
3193 },
3194 LinkedHttpRouteOwner {
3195 module_name: "auth-password".to_owned(),
3196 public_prefixes: &["/v1/auth/password/"],
3197 },
3198 LinkedHttpRouteOwner {
3199 module_name: "platform-story".to_owned(),
3200 public_prefixes: &["/admin/runtime/stories"],
3201 },
3202 ]
3203 );
3204 }
3205
3206 #[test]
3207 fn linked_http_bindings_are_declared_in_manifests() {
3208 for module in linked_http_modules() {
3209 let http = module
3210 .linked_http
3211 .expect("linked HTTP module should carry HTTP contribution");
3212 assert!(
3213 !module.manifest.http_routes.is_empty(),
3214 "linked HTTP module `{}` must declare ModuleManifest::http_routes",
3215 module.manifest.name
3216 );
3217 for route in &module.manifest.http_routes {
3218 assert!(
3219 http.public_prefixes
3220 .iter()
3221 .any(|prefix| route.path.starts_with(prefix)),
3222 "linked HTTP module `{}` declares manifest route `{}` outside its public prefixes",
3223 module.manifest.name,
3224 route.path
3225 );
3226 }
3227 }
3228 }
3229
3230 #[test]
3231 fn linked_http_modules_are_registered_modules() {
3232 let manifests = module_manifests();
3233
3234 for module in linked_http_modules() {
3235 let registered_manifest = manifests
3236 .iter()
3237 .find(|manifest| manifest.name == module.manifest.name)
3238 .unwrap_or_else(|| {
3239 panic!(
3240 "linked HTTP module `{}` is missing from module_manifests",
3241 module.manifest.name
3242 )
3243 });
3244 assert_eq!(
3245 registered_manifest, &module.manifest,
3246 "linked HTTP module `{}` must use the registered ModuleManifest",
3247 module.manifest.name
3248 );
3249 }
3250 }
3251
3252 #[test]
3253 fn linked_http_routes_include_story_module_routes() {
3254 let document = merge_linked_http(platform_http::OpenApiRouter::new()).to_openapi();
3255 let value = serde_json::to_value(document).expect("OpenAPI document should serialize");
3256 let paths = value["paths"].as_object().expect("OpenAPI paths object");
3257
3258 assert!(paths.contains_key("/admin/runtime/stories"));
3259 assert!(paths.contains_key("/admin/runtime/stories/{correlation_id}"));
3260 assert!(paths.contains_key("/admin/runtime/stories/{correlation_id}/heatmap"));
3261 assert!(paths.contains_key("/admin/runtime/stories/{correlation_id}/technical-operations"));
3262 }
3263
3264 #[test]
3265 fn platform_story_manifest_declares_story_console_surface() {
3266 let manifest = module_manifests()
3267 .into_iter()
3268 .find(|manifest| manifest.name == "platform-story")
3269 .expect("platform-story manifest should be registered");
3270 let console_surface_contract: Value = serde_json::from_str(include_str!(
3271 "../../../modules/story/console/console-surface.json"
3272 ))
3273 .expect("story console surface contract should be valid json");
3274
3275 assert_eq!(manifest.admin, None);
3276 assert_eq!(manifest.console.len(), 1);
3277 let surface = &manifest.console[0];
3278 let surface_json =
3279 serde_json::to_value(surface).expect("platform-story console surface should serialize");
3280
3281 assert_eq!(
3282 manifest.capabilities,
3283 required_capabilities_from_contract(&console_surface_contract)
3284 );
3285 assert_eq!(manifest.name, console_surface_contract["id"]);
3286 assert_eq!(surface.name, console_surface_contract["surfaceName"]);
3287 assert_eq!(surface.label, console_surface_contract["label"]);
3288 assert_eq!(surface.area, ConsoleArea::Runtime);
3289 assert_eq!(surface_json["area"], console_surface_contract["area"]);
3290 assert_eq!(surface.route, console_surface_contract["route"]);
3291 assert_eq!(
3292 surface.package.name,
3293 console_surface_contract["packageName"]
3294 );
3295 assert_eq!(
3296 surface.package.export,
3297 console_surface_contract["exportName"]
3298 );
3299 assert_eq!(surface_json["icon"], console_surface_contract["icon"]);
3300 assert_eq!(surface.navigation, None);
3301 assert!(console_surface_contract.get("navigation").is_none());
3302 assert_eq!(
3303 surface.required_capabilities,
3304 required_capabilities_from_contract(&console_surface_contract)
3305 );
3306
3307 let lints = lint_module_manifest(ModuleSource::Linked, &manifest);
3308 assert!(
3309 lints
3310 .iter()
3311 .all(|lint| lint.severity == ModuleManifestLintSeverity::Ok),
3312 "platform-story manifest should not have warning/error lints: {lints:?}"
3313 );
3314 }
3315
3316 fn required_capabilities_from_contract(contract: &Value) -> Vec<String> {
3317 contract["requiredCapabilities"]
3318 .as_array()
3319 .expect("requiredCapabilities should be an array")
3320 .iter()
3321 .map(|capability| {
3322 capability
3323 .as_str()
3324 .expect("requiredCapabilities should contain strings")
3325 .to_owned()
3326 })
3327 .collect()
3328 }
3329
3330 #[tokio::test]
3331 async fn lifecycle_activation_enqueue_creates_function_run() {
3332 let Some(db) = TestDatabase::create().await else {
3333 return;
3334 };
3335 apply_runtime_stack_migrations(&db).await;
3336
3337 let mut ctx = AppContext::new(
3338 test_config(&db),
3339 db.pool.clone(),
3340 Arc::new(LoggingEventPublisher),
3341 );
3342 ctx.ids = Arc::new(SequentialIdGenerator::default());
3343 let modules = vec![
3344 test_lifecycle_module(lifecycle_activation_job(true, json!({ "warm": "cache" })))
3345 .into(),
3346 ];
3347 let registry = registry_with_lifecycle_function(7);
3348
3349 let run_ids = enqueue_lifecycle_activation_jobs(&ctx, &modules, ®istry)
3350 .await
3351 .expect("lifecycle activation job should enqueue");
3352
3353 assert_eq!(run_ids.len(), 1);
3354 let row = sqlx::query_as::<_, (String, Value, i32, String, Value)>(
3355 r#"
3356 select function_name, input_json, max_attempts, correlation_id, actor
3357 from runtime.function_runs
3358 where id = $1
3359 "#,
3360 )
3361 .bind(&run_ids[0])
3362 .fetch_one(&db.pool)
3363 .await
3364 .expect("function run should exist");
3365
3366 assert_eq!(row.0, LIFECYCLE_FUNCTION_NAME);
3367 assert_eq!(row.1["warm"], "cache");
3368 assert_eq!(
3369 row.1["_lenso_runtime"]["correlation_id"],
3370 "corr_lifecycle_1"
3371 );
3372 assert_eq!(
3373 row.1["_lenso_runtime"]["causation_id"],
3374 "module_lifecycle:test-module:warm cache"
3375 );
3376 assert_eq!(row.2, 7);
3377 assert_eq!(row.3, "corr_lifecycle_1");
3378 assert_eq!(row.4["kind"], "service");
3379 assert_eq!(row.4["service_id"], "worker");
3380 assert_eq!(row.4["scopes"][0], "runtime.functions.enqueue");
3381
3382 db.cleanup().await;
3383 }
3384
3385 #[test]
3386 fn lifecycle_activation_validation_rejects_required_missing_function() {
3387 let modules =
3388 vec![test_lifecycle_module(lifecycle_activation_job(true, Value::Null)).into()];
3389 let registry = FunctionRegistry::default();
3390
3391 let error = validate_lifecycle_activation_jobs(&modules, ®istry)
3392 .expect_err("required missing activation function should fail validation");
3393
3394 assert_eq!(error.code, ErrorCode::Validation);
3395 assert_eq!(
3396 error.details[0].field.as_deref(),
3397 Some("module.test-module.lifecycle.activation_jobs.warm cache")
3398 );
3399 assert!(
3400 error.details[0].reason.contains("missing function"),
3401 "validation detail should name the missing registry function"
3402 );
3403 }
3404
3405 #[test]
3406 fn lifecycle_activation_validation_rejects_required_startup_check_missing_function() {
3407 let modules = vec![test_lifecycle_module_with_lifecycle(
3408 LifecycleSurface {
3409 startup_checks: vec![LifecycleStartupCheckDeclaration {
3410 name: "function registered".to_owned(),
3411 required: true,
3412 check: LifecycleStartupCheckKind::FunctionRegistered {
3413 function_name: LIFECYCLE_FUNCTION_NAME.to_owned(),
3414 },
3415 }],
3416 activation_jobs: Vec::new(),
3417 },
3418 true,
3419 Vec::new(),
3420 )];
3421 let registry = FunctionRegistry::default();
3422
3423 let error = validate_lifecycle_activation_jobs(&modules, ®istry)
3424 .expect_err("required startup check should fail when function is missing");
3425
3426 assert_eq!(error.code, ErrorCode::Validation);
3427 assert_eq!(
3428 error.details[0].field.as_deref(),
3429 Some("module.test-module.lifecycle.startup_checks.function registered")
3430 );
3431 assert!(
3432 error.details[0].reason.contains("missing function"),
3433 "validation detail should name the missing registry function"
3434 );
3435 }
3436
3437 #[test]
3438 fn lifecycle_activation_validation_rejects_required_startup_check_function_not_declared() {
3439 let modules = vec![test_lifecycle_module_with_lifecycle(
3440 LifecycleSurface {
3441 startup_checks: vec![LifecycleStartupCheckDeclaration {
3442 name: "function registered".to_owned(),
3443 required: true,
3444 check: LifecycleStartupCheckKind::FunctionRegistered {
3445 function_name: LIFECYCLE_FUNCTION_NAME.to_owned(),
3446 },
3447 }],
3448 activation_jobs: Vec::new(),
3449 },
3450 false,
3451 Vec::new(),
3452 )];
3453 let registry = registry_with_lifecycle_function(3);
3454
3455 let error = validate_lifecycle_activation_jobs(&modules, ®istry)
3456 .expect_err("required startup check should fail when manifest does not declare it");
3457
3458 assert_eq!(error.code, ErrorCode::Validation);
3459 assert_eq!(
3460 error.details[0].field.as_deref(),
3461 Some("module.test-module.lifecycle.startup_checks.function registered")
3462 );
3463 assert!(
3464 error.details[0].reason.contains("not declared"),
3465 "validation detail should name the missing module runtime declaration"
3466 );
3467 }
3468
3469 #[test]
3470 fn lifecycle_activation_validation_rejects_required_startup_check_missing_capability() {
3471 let modules = vec![test_lifecycle_module_with_lifecycle(
3472 LifecycleSurface {
3473 startup_checks: vec![LifecycleStartupCheckDeclaration {
3474 name: "capability declared".to_owned(),
3475 required: true,
3476 check: LifecycleStartupCheckKind::CapabilityDeclared {
3477 capability: "test.cache.warm".to_owned(),
3478 },
3479 }],
3480 activation_jobs: Vec::new(),
3481 },
3482 false,
3483 Vec::new(),
3484 )];
3485 let registry = FunctionRegistry::default();
3486
3487 let error = validate_lifecycle_activation_jobs(&modules, ®istry)
3488 .expect_err("required startup check should fail when capability is missing");
3489
3490 assert_eq!(error.code, ErrorCode::Validation);
3491 assert_eq!(
3492 error.details[0].field.as_deref(),
3493 Some("module.test-module.lifecycle.startup_checks.capability declared")
3494 );
3495 assert!(
3496 error.details[0].reason.contains("missing capability"),
3497 "validation detail should name the missing capability"
3498 );
3499 }
3500
3501 #[test]
3502 fn lifecycle_activation_optional_startup_checks_do_not_fail_validation() {
3503 let modules = vec![test_lifecycle_module_with_lifecycle(
3504 LifecycleSurface {
3505 startup_checks: vec![
3506 LifecycleStartupCheckDeclaration {
3507 name: "optional function".to_owned(),
3508 required: false,
3509 check: LifecycleStartupCheckKind::FunctionRegistered {
3510 function_name: LIFECYCLE_FUNCTION_NAME.to_owned(),
3511 },
3512 },
3513 LifecycleStartupCheckDeclaration {
3514 name: "optional capability".to_owned(),
3515 required: false,
3516 check: LifecycleStartupCheckKind::CapabilityDeclared {
3517 capability: "test.cache.warm".to_owned(),
3518 },
3519 },
3520 ],
3521 activation_jobs: Vec::new(),
3522 },
3523 false,
3524 Vec::new(),
3525 )];
3526 let registry = FunctionRegistry::default();
3527
3528 validate_lifecycle_activation_jobs(&modules, ®istry)
3529 .expect("optional startup checks should not fail validation");
3530 }
3531
3532 #[test]
3533 fn lifecycle_activation_validation_rejects_required_job_not_declared_by_module() {
3534 let modules = vec![
3535 test_lifecycle_module(lifecycle_activation_job(true, Value::Null))
3536 .without_runtime_declaration()
3537 .into(),
3538 ];
3539 let registry = registry_with_lifecycle_function(3);
3540
3541 let error = validate_lifecycle_activation_jobs(&modules, ®istry)
3542 .expect_err("required activation job should fail when manifest does not declare it");
3543
3544 assert_eq!(error.code, ErrorCode::Validation);
3545 assert_eq!(
3546 error.details[0].field.as_deref(),
3547 Some("module.test-module.lifecycle.activation_jobs.warm cache")
3548 );
3549 assert!(
3550 error.details[0].reason.contains("not declared"),
3551 "validation detail should name the missing module runtime declaration"
3552 );
3553 }
3554
3555 #[tokio::test]
3556 async fn optional_missing_lifecycle_activation_is_skipped() {
3557 let db = platform_core::DbPool::connect_lazy("postgres://localhost/lenso_test")
3558 .expect("lazy pool should build");
3559 let ctx = AppContext::new(
3560 test_config_with_database_url("postgres://localhost/lenso_test"),
3561 db,
3562 Arc::new(LoggingEventPublisher),
3563 );
3564 let modules =
3565 vec![test_lifecycle_module(lifecycle_activation_job(false, Value::Null)).into()];
3566 let registry = FunctionRegistry::default();
3567
3568 let run_ids = enqueue_lifecycle_activation_jobs(&ctx, &modules, ®istry)
3569 .await
3570 .expect("optional missing activation function should be skipped");
3571
3572 assert!(run_ids.is_empty());
3573 }
3574
3575 #[tokio::test]
3576 async fn lifecycle_activation_optional_job_not_declared_is_skipped() {
3577 let db = platform_core::DbPool::connect_lazy("postgres://localhost/lenso_test")
3578 .expect("lazy pool should build");
3579 let ctx = AppContext::new(
3580 test_config_with_database_url("postgres://localhost/lenso_test"),
3581 db,
3582 Arc::new(LoggingEventPublisher),
3583 );
3584 let modules = vec![
3585 test_lifecycle_module(lifecycle_activation_job(false, Value::Null))
3586 .without_runtime_declaration()
3587 .into(),
3588 ];
3589 let registry = registry_with_lifecycle_function(3);
3590
3591 let run_ids = enqueue_lifecycle_activation_jobs(&ctx, &modules, ®istry)
3592 .await
3593 .expect("optional undeclared activation function should be skipped");
3594
3595 assert!(run_ids.is_empty());
3596 }
3597
3598 #[tokio::test]
3599 async fn lifecycle_activation_optional_enqueue_failure_is_skipped() {
3600 let db = PgPoolOptions::new()
3601 .max_connections(1)
3602 .acquire_timeout(Duration::from_millis(50))
3603 .connect_lazy_with(
3604 PgConnectOptions::new()
3605 .host("127.0.0.1")
3606 .port(1)
3607 .username("postgres")
3608 .database("lenso_test"),
3609 );
3610 let ctx = AppContext::new(
3611 test_config_with_database_url("postgres://localhost:1/lenso_test"),
3612 db,
3613 Arc::new(LoggingEventPublisher),
3614 );
3615 let modules =
3616 vec![test_lifecycle_module(lifecycle_activation_job(false, Value::Null)).into()];
3617 let registry = registry_with_lifecycle_function(3);
3618
3619 let run_ids = enqueue_lifecycle_activation_jobs(&ctx, &modules, ®istry)
3620 .await
3621 .expect("optional enqueue failure should be skipped");
3622
3623 assert!(run_ids.is_empty());
3624 }
3625
3626 #[test]
3627 fn lifecycle_activation_max_attempts_conversion_saturates() {
3628 assert_eq!(runtime_max_attempts_for_enqueue(7), 7);
3629 assert_eq!(runtime_max_attempts_for_enqueue(u32::MAX), i32::MAX);
3630 }
3631
3632 const LIFECYCLE_FUNCTION_NAME: &str = "test.warm_cache.v1";
3633
3634 #[derive(Debug)]
3635 struct NoopFunctionHandler;
3636
3637 #[async_trait]
3638 impl FunctionHandler for NoopFunctionHandler {
3639 async fn call(
3640 &self,
3641 _ctx: ExecutionContext,
3642 _input: Value,
3643 ) -> platform_core::AppResult<Value> {
3644 Ok(Value::Null)
3645 }
3646 }
3647
3648 fn lifecycle_activation_job(required: bool, input: Value) -> LifecycleActivationJobDeclaration {
3649 LifecycleActivationJobDeclaration {
3650 name: "warm cache".to_owned(),
3651 function_name: LIFECYCLE_FUNCTION_NAME.to_owned(),
3652 run_policy: LifecycleActivationRunPolicy::EveryStartup,
3653 input,
3654 required,
3655 }
3656 }
3657
3658 struct TestLifecycleModuleBuilder {
3659 lifecycle: LifecycleSurface,
3660 declare_runtime_function: bool,
3661 capabilities: Vec<String>,
3662 }
3663
3664 impl TestLifecycleModuleBuilder {
3665 fn without_runtime_declaration(mut self) -> Self {
3666 self.declare_runtime_function = false;
3667 self
3668 }
3669 }
3670
3671 impl From<TestLifecycleModuleBuilder> for Module {
3672 fn from(builder: TestLifecycleModuleBuilder) -> Self {
3673 let mut manifest = ModuleManifest::builder("test-module").lifecycle(builder.lifecycle);
3674 if builder.declare_runtime_function {
3675 manifest = manifest.runtime(RuntimeSurface {
3676 functions: vec![RuntimeFunctionDeclaration {
3677 name: LIFECYCLE_FUNCTION_NAME.to_owned(),
3678 version: 1,
3679 queue: "test".to_owned(),
3680 input_schema: None,
3681 retry_policy: None,
3682 }],
3683 });
3684 }
3685 if !builder.capabilities.is_empty() {
3686 manifest = manifest.capabilities(builder.capabilities);
3687 }
3688 Module::linked(manifest.build(), LinkedBinding::builder().build())
3689 }
3690 }
3691
3692 fn test_lifecycle_module(job: LifecycleActivationJobDeclaration) -> TestLifecycleModuleBuilder {
3693 TestLifecycleModuleBuilder {
3694 lifecycle: LifecycleSurface {
3695 startup_checks: Vec::new(),
3696 activation_jobs: vec![job],
3697 },
3698 declare_runtime_function: true,
3699 capabilities: Vec::new(),
3700 }
3701 }
3702
3703 fn test_lifecycle_module_with_lifecycle(
3704 lifecycle: LifecycleSurface,
3705 declare_runtime_function: bool,
3706 capabilities: Vec<String>,
3707 ) -> Module {
3708 TestLifecycleModuleBuilder {
3709 lifecycle,
3710 declare_runtime_function,
3711 capabilities,
3712 }
3713 .into()
3714 }
3715
3716 fn registry_with_lifecycle_function(max_attempts: u32) -> FunctionRegistry {
3717 let mut registry = FunctionRegistry::default();
3718 registry.register(FunctionDefinition {
3719 name: LIFECYCLE_FUNCTION_NAME.to_owned(),
3720 version: 1,
3721 queue: "test".to_owned(),
3722 retry_policy: RetryPolicy::fixed(max_attempts, Duration::ZERO),
3723 handler: Arc::new(NoopFunctionHandler),
3724 });
3725 registry
3726 }
3727
3728 #[test]
3729 fn remote_module_service_specs_parse() {
3730 let specs = parse_remote_module_service_specs(&serde_json::json!({
3731 "modules": [
3732 {
3733 "moduleName": "crm",
3734 "services": [
3735 {
3736 "name": "crm-api",
3737 "command": "pnpm dev",
3738 "cwd": "../crm",
3739 "readyUrl": "http://127.0.0.1:4100/lenso/module/v1/manifest",
3740 "readyTimeoutMs": 12000,
3741 "autoStart": true
3742 }
3743 ]
3744 }
3745 ],
3746 "version": 1
3747 }))
3748 .expect("service specs parse");
3749
3750 assert_eq!(specs.len(), 1);
3751 assert_eq!(specs[0].module_name, "crm");
3752 assert_eq!(specs[0].service_name, "crm-api");
3753 assert_eq!(specs[0].ready_timeout_ms, 12000);
3754 }
3755
3756 #[test]
3757 fn remote_module_service_state_path_sanitizes_names() {
3758 let spec = RemoteModuleServiceSpec {
3759 module_name: "CRM Module".to_owned(),
3760 service_name: "API Worker!".to_owned(),
3761 command: "pnpm dev".to_owned(),
3762 cwd: None,
3763 ready_url: "http://127.0.0.1:4100/lenso/module/v1/manifest".to_owned(),
3764 ready_timeout_ms: 12000,
3765 auto_start: true,
3766 };
3767
3768 let path = remote_module_service_state_path(Path::new(".lenso"), &spec, "lock");
3769
3770 assert_eq!(
3771 path,
3772 PathBuf::from(".lenso/remote-crm-module-api-worker.lock")
3773 );
3774 }
3775
3776 #[test]
3777 fn remote_module_service_lock_is_exclusive_and_released() {
3778 let unique = std::time::SystemTime::now()
3779 .duration_since(std::time::UNIX_EPOCH)
3780 .expect("system time should be after Unix epoch")
3781 .as_nanos();
3782 let dir = std::env::temp_dir().join(format!(
3783 "lenso-bootstrap-service-lock-{}-{unique}",
3784 std::process::id()
3785 ));
3786 let lock_file_path = dir.join("service.lock");
3787 let pid_file_path = dir.join("service.pid");
3788
3789 let _ = std::fs::remove_dir_all(&dir);
3790 create_remote_module_service_lock(&lock_file_path)
3791 .expect("first lock claim should create the lock");
3792 let second_claim = create_remote_module_service_lock(&lock_file_path)
3793 .expect_err("second lock claim should fail while the file exists");
3794 assert_eq!(second_claim.kind(), std::io::ErrorKind::AlreadyExists);
3795 std::fs::write(&pid_file_path, "123\n").expect("pid file should write");
3796
3797 release_remote_module_service_state(&lock_file_path, &pid_file_path);
3798
3799 assert!(!lock_file_path.exists());
3800 assert!(!pid_file_path.exists());
3801 let _ = std::fs::remove_dir_all(&dir);
3802 }
3803
3804 const TEST_HOST_MIGRATIONS: &[Migration] = &[Migration {
3805 name: "billing/0001_init",
3806 sql: "select 1;",
3807 }];
3808
3809 fn test_host_manifest() -> ModuleManifest {
3810 ModuleManifest::builder("billing").build()
3811 }
3812
3813 fn test_host_linked_module() -> HostLinkedModule {
3814 HostLinkedModule::manifest_only("billing", test_host_manifest, TEST_HOST_MIGRATIONS)
3815 }
3816
3817 fn test_config(db: &TestDatabase) -> AppConfig {
3818 test_config_with_database_url(db.url.clone())
3819 }
3820
3821 fn test_config_with_database_url(database_url: impl Into<String>) -> AppConfig {
3822 AppConfig {
3823 service: ServiceConfig::default(),
3824 database: DatabaseConfig {
3825 url: database_url.into(),
3826 max_connections: 5,
3827 },
3828 redis: RedisConfig::default(),
3829 http: HttpConfig::default(),
3830 telemetry: TelemetryConfig::default(),
3831 auth: AuthConfig::default(),
3832 console: Default::default(),
3833 module_sources: ModuleSourcesConfig::default(),
3834 modules: BTreeMap::new(),
3835 }
3836 }
3837
3838 async fn apply_runtime_stack_migrations(db: &TestDatabase) {
3839 let migrations = PLATFORM_MIGRATIONS
3840 .iter()
3841 .chain(RUNTIME_MIGRATIONS)
3842 .copied()
3843 .collect::<Vec<_>>();
3844 apply_migrations(&db.pool, &migrations)
3845 .await
3846 .expect("platform and runtime migrations should apply");
3847 }
3848}