1use std::{path::PathBuf, sync::Arc};
4
5use aion::{
6 ActivityDispatcher, EngineBuilder, RuntimeHandle, SignalRouter, signal::ConcreteSignalRouter,
7};
8use aion_store::{EventStore, NamespaceStore, OutboxStore};
9#[cfg(feature = "libsql-backend")]
10use aion_store_libsql::LibSqlStore;
11
12use crate::dev_ui::{ActivityMockRegistry, DevMockingDispatcher};
13
14#[cfg(feature = "auth")]
15use crate::auth::JwksCache;
16use crate::{
17 config::{RuntimeConfig, ServerConfig, StoreBackend, StoreConfig},
18 error::ServerError,
19 namespace::{NamespaceGuard, NamespaceMinter, resolver::NamespaceResolver},
20 observability::{
21 Metrics, health::HealthState, instrumented_store::InstrumentedEventStore,
22 metrics::MetricsError,
23 },
24 shutdown::DrainState,
25 worker::{
26 ConnectedWorkerRegistry, HeartbeatTracker, PendingActivities, WorkerActivityDispatcher,
27 },
28};
29
30#[derive(Clone)]
32pub struct ServerState {
33 inner: Arc<ServerStateInner>,
34}
35
36struct ServerStateInner {
37 namespace_guard: NamespaceGuard,
38 runtime: RuntimeConfig,
39 worker_registry: ConnectedWorkerRegistry,
40 pending_activities: PendingActivities,
41 heartbeat_tracker: HeartbeatTracker,
42 drain_state: DrainState,
43 metrics: Option<Metrics>,
44 health: Option<HealthState>,
45 activity_mock_registry: Option<ActivityMockRegistry>,
48 outbox_store: Option<Arc<dyn OutboxStore>>,
52 namespace_store: Arc<dyn NamespaceStore>,
61 outbox_wake: Arc<tokio::sync::Notify>,
68 cluster_publisher: crate::cluster_publisher::ClusterEventPublisher,
72 transcript_publisher: crate::activity_publisher::ActivityEventPublisher,
81 attempt_owners: crate::worker::AttemptOwnerIndex,
88 cluster_self_node: Option<String>,
93 #[cfg(feature = "haematite-backend")]
98 cluster_responder: Option<aion_store_haematite::ClusterResponder>,
99 #[cfg(feature = "haematite-backend")]
102 cluster_store: Option<Arc<aion_store_haematite::HaematiteStore>>,
103 #[cfg(feature = "haematite-backend")]
106 watched_peers: Vec<crate::cluster::WatchedPeer>,
107 #[cfg(feature = "haematite-backend")]
112 shard_directory: Option<Arc<crate::routing::StaticShardDirectory>>,
113 #[cfg(feature = "haematite-backend")]
117 request_forwarder: Option<Arc<dyn crate::routing::RequestForwarder>>,
118 #[cfg(feature = "auth")]
119 jwks_cache: Option<JwksCache>,
120}
121
122impl ServerState {
123 const FALLBACK_CLUSTER_BROADCAST_CAPACITY: std::num::NonZeroUsize =
132 match std::num::NonZeroUsize::new(64) {
133 Some(value) => value,
134 None => std::num::NonZeroUsize::MIN,
135 };
136
137 pub async fn build(config: ServerConfig) -> Result<Self, ServerError> {
144 let (store_config, runtime) = config.into_parts();
145 let connected = connect_store(store_config).await?;
146 Self::build_with_connected_store(connected, runtime).await
147 }
148
149 pub async fn build_with_store<S>(store: S, runtime: RuntimeConfig) -> Result<Self, ServerError>
155 where
156 S: EventStore + NamespaceStore,
157 {
158 let leaf = Arc::new(store);
162 let namespace_store: Arc<dyn NamespaceStore> = leaf.clone();
163 Self::build_with_connected_store(
164 ConnectedStore::local(leaf, None, namespace_store),
165 runtime,
166 )
167 .await
168 }
169
170 async fn build_with_connected_store(
171 connected: ConnectedStore,
172 runtime: RuntimeConfig,
173 ) -> Result<Self, ServerError> {
174 let outbox_store = connected.outbox_store;
175 let bootstrap_coordinator = connected.bootstrap_coordinator;
176 #[cfg(feature = "haematite-backend")]
177 let cluster_responder = connected.cluster_responder;
178 #[cfg(feature = "haematite-backend")]
179 let cluster_store = connected.cluster_store;
180 #[cfg(feature = "haematite-backend")]
181 let watched_peers = connected.watched_peers;
182 #[cfg(feature = "haematite-backend")]
185 let cluster_self_node = connected.self_node_id.clone();
186 #[cfg(not(feature = "haematite-backend"))]
187 let cluster_self_node: Option<String> = None;
188 #[cfg(feature = "haematite-backend")]
192 let RoutingState {
193 shard_directory,
194 request_forwarder,
195 } = build_routing_state(
196 cluster_store.as_ref(),
197 connected.directory_peers,
198 connected.self_node_id,
199 );
200 let (event_broadcast_capacity, query_timeout) = required_engine_seams(&runtime)?;
201 let (cluster_publisher, transcript_publisher) =
202 build_real_time_publishers(&runtime, connected.observability_store)?;
203 let metrics = Metrics::new().map_err(|error| metrics_config_error(&error))?;
204 let outbox_wake = Arc::new(tokio::sync::Notify::new());
209 let instrumented_store = Arc::new(
210 InstrumentedEventStore::new(
211 connected.event_store,
212 metrics.clone(),
213 runtime.default_namespace.clone(),
214 )
215 .with_outbox_wake(Arc::clone(&outbox_wake)),
216 );
217 let exported_metrics = runtime.metrics.enabled.then_some(metrics.clone());
218 let worker_registry = ConnectedWorkerRegistry::default()
220 .with_cluster_publisher(cluster_publisher.clone())
221 .with_namespace_minting(connected.namespace_store.clone(), runtime.auto_create);
222 let pending_activities = PendingActivities::default();
223 let heartbeat_tracker = HeartbeatTracker::new(runtime.worker.heartbeat_window);
224 let drain_state = DrainState::default();
225 let (dispatcher, attempt_owners) = build_bridge_dispatcher(
226 &runtime,
227 &worker_registry,
228 &pending_activities,
229 &heartbeat_tracker,
230 &drain_state,
231 );
232 let (activity_dispatcher, activity_mock_registry) =
233 decorate_activity_dispatcher(dispatcher, runtime.dev.enabled);
234
235 let engine = build_engine(EngineAssembly {
236 instrumented_store: &instrumented_store,
237 event_broadcast_capacity,
238 query_timeout,
239 activity_dispatcher,
240 active_registry: Arc::new(aion::Registry::default()),
241 bootstrap_coordinator,
242 runtime: &runtime,
243 })
244 .await?;
245 let engine = Arc::new(engine);
246 install_outbox_delivery(&pending_activities, &engine, runtime.outbox.enabled);
247 let resolver = NamespaceResolver::from_config(runtime.namespace.clone(), engine);
248 #[cfg(feature = "auth")]
249 let jwks_cache = build_jwks_cache(&runtime).await?;
250 Ok(Self {
251 inner: Arc::new(ServerStateInner {
252 namespace_guard: NamespaceGuard::new(resolver),
253 runtime,
254 worker_registry,
255 pending_activities,
256 heartbeat_tracker,
257 drain_state,
258 metrics: exported_metrics,
259 health: Some(HealthState::new(instrumented_store, true)),
260 activity_mock_registry,
261 outbox_store,
262 namespace_store: connected.namespace_store,
263 outbox_wake,
264 cluster_publisher,
265 transcript_publisher,
266 attempt_owners,
267 cluster_self_node,
268 #[cfg(feature = "haematite-backend")]
269 cluster_responder,
270 #[cfg(feature = "haematite-backend")]
271 cluster_store,
272 #[cfg(feature = "haematite-backend")]
273 watched_peers,
274 #[cfg(feature = "haematite-backend")]
275 shard_directory,
276 #[cfg(feature = "haematite-backend")]
277 request_forwarder,
278 #[cfg(feature = "auth")]
279 jwks_cache,
280 }),
281 })
282 }
283
284 #[must_use]
286 pub fn from_parts(namespace_resolver: NamespaceResolver, runtime: RuntimeConfig) -> Self {
287 Self::from_parts_with_namespace_store(
292 namespace_resolver,
293 runtime,
294 Arc::new(aion_store::InMemoryStore::default()),
295 )
296 }
297
298 #[must_use]
306 pub fn from_parts_with_namespace_store(
307 namespace_resolver: NamespaceResolver,
308 runtime: RuntimeConfig,
309 namespace_store: Arc<dyn NamespaceStore>,
310 ) -> Self {
311 let heartbeat_tracker = HeartbeatTracker::new(runtime.worker.heartbeat_window);
312 let bounds = transcript_bounds(&runtime);
316 Self {
317 inner: Arc::new(ServerStateInner {
318 namespace_guard: NamespaceGuard::new(namespace_resolver),
319 runtime,
320 worker_registry: ConnectedWorkerRegistry::default(),
321 pending_activities: PendingActivities::default(),
322 heartbeat_tracker,
323 drain_state: DrainState::default(),
324 metrics: None,
325 health: None,
326 activity_mock_registry: None,
327 outbox_store: None,
328 namespace_store,
329 outbox_wake: Arc::new(tokio::sync::Notify::new()),
330 cluster_publisher: crate::cluster_publisher::ClusterEventPublisher::new(
331 Self::FALLBACK_CLUSTER_BROADCAST_CAPACITY,
332 ),
333 transcript_publisher: build_transcript_publisher(
337 None,
338 Self::FALLBACK_CLUSTER_BROADCAST_CAPACITY,
339 bounds,
340 ),
341 attempt_owners: crate::worker::AttemptOwnerIndex::new(),
342 cluster_self_node: None,
343 #[cfg(feature = "haematite-backend")]
344 cluster_responder: None,
345 #[cfg(feature = "haematite-backend")]
346 cluster_store: None,
347 #[cfg(feature = "haematite-backend")]
348 watched_peers: Vec::new(),
349 #[cfg(feature = "haematite-backend")]
350 shard_directory: None,
351 #[cfg(feature = "haematite-backend")]
352 request_forwarder: None,
353 #[cfg(feature = "auth")]
354 jwks_cache: None,
355 }),
356 }
357 }
358
359 #[cfg(feature = "auth")]
368 #[must_use]
369 pub fn from_parts_with_namespace_store_and_jwks(
370 namespace_resolver: NamespaceResolver,
371 runtime: RuntimeConfig,
372 namespace_store: Arc<dyn NamespaceStore>,
373 jwks_cache: JwksCache,
374 ) -> Self {
375 let heartbeat_tracker = HeartbeatTracker::new(runtime.worker.heartbeat_window);
376 let bounds = transcript_bounds(&runtime);
380 Self {
381 inner: Arc::new(ServerStateInner {
382 namespace_guard: NamespaceGuard::new(namespace_resolver),
383 runtime,
384 worker_registry: ConnectedWorkerRegistry::default(),
385 pending_activities: PendingActivities::default(),
386 heartbeat_tracker,
387 drain_state: DrainState::default(),
388 metrics: None,
389 health: None,
390 activity_mock_registry: None,
391 outbox_store: None,
392 namespace_store,
393 outbox_wake: Arc::new(tokio::sync::Notify::new()),
394 cluster_publisher: crate::cluster_publisher::ClusterEventPublisher::new(
395 Self::FALLBACK_CLUSTER_BROADCAST_CAPACITY,
396 ),
397 transcript_publisher: build_transcript_publisher(
401 None,
402 Self::FALLBACK_CLUSTER_BROADCAST_CAPACITY,
403 bounds,
404 ),
405 attempt_owners: crate::worker::AttemptOwnerIndex::new(),
406 cluster_self_node: None,
407 #[cfg(feature = "haematite-backend")]
408 cluster_responder: None,
409 #[cfg(feature = "haematite-backend")]
410 cluster_store: None,
411 #[cfg(feature = "haematite-backend")]
412 watched_peers: Vec::new(),
413 #[cfg(feature = "haematite-backend")]
414 shard_directory: None,
415 #[cfg(feature = "haematite-backend")]
416 request_forwarder: None,
417 jwks_cache: Some(jwks_cache),
418 }),
419 }
420 }
421
422 #[cfg(feature = "auth")]
428 #[must_use]
429 pub fn from_parts_with_jwks(
430 namespace_resolver: NamespaceResolver,
431 runtime: RuntimeConfig,
432 jwks_cache: JwksCache,
433 ) -> Self {
434 let heartbeat_tracker = HeartbeatTracker::new(runtime.worker.heartbeat_window);
435 let bounds = transcript_bounds(&runtime);
439 Self {
440 inner: Arc::new(ServerStateInner {
441 namespace_guard: NamespaceGuard::new(namespace_resolver),
442 runtime,
443 worker_registry: ConnectedWorkerRegistry::default(),
444 pending_activities: PendingActivities::default(),
445 heartbeat_tracker,
446 drain_state: DrainState::default(),
447 metrics: None,
448 health: None,
449 activity_mock_registry: None,
450 outbox_store: None,
451 namespace_store: Arc::new(aion_store::InMemoryStore::default()),
456 outbox_wake: Arc::new(tokio::sync::Notify::new()),
457 cluster_publisher: crate::cluster_publisher::ClusterEventPublisher::new(
458 Self::FALLBACK_CLUSTER_BROADCAST_CAPACITY,
459 ),
460 transcript_publisher: build_transcript_publisher(
464 None,
465 Self::FALLBACK_CLUSTER_BROADCAST_CAPACITY,
466 bounds,
467 ),
468 attempt_owners: crate::worker::AttemptOwnerIndex::new(),
469 cluster_self_node: None,
470 #[cfg(feature = "haematite-backend")]
471 cluster_responder: None,
472 #[cfg(feature = "haematite-backend")]
473 cluster_store: None,
474 #[cfg(feature = "haematite-backend")]
475 watched_peers: Vec::new(),
476 #[cfg(feature = "haematite-backend")]
477 shard_directory: None,
478 #[cfg(feature = "haematite-backend")]
479 request_forwarder: None,
480 jwks_cache: Some(jwks_cache),
481 }),
482 }
483 }
484
485 #[must_use]
487 pub fn from_parts_with_registry(
488 namespace_resolver: NamespaceResolver,
489 runtime: RuntimeConfig,
490 worker_registry: ConnectedWorkerRegistry,
491 ) -> Self {
492 let heartbeat_tracker = HeartbeatTracker::new(runtime.worker.heartbeat_window);
493 let bounds = transcript_bounds(&runtime);
497 Self {
498 inner: Arc::new(ServerStateInner {
499 namespace_guard: NamespaceGuard::new(namespace_resolver),
500 runtime,
501 worker_registry,
502 pending_activities: PendingActivities::default(),
503 heartbeat_tracker,
504 drain_state: DrainState::default(),
505 metrics: None,
506 health: None,
507 activity_mock_registry: None,
508 outbox_store: None,
509 namespace_store: Arc::new(aion_store::InMemoryStore::default()),
514 outbox_wake: Arc::new(tokio::sync::Notify::new()),
515 cluster_publisher: crate::cluster_publisher::ClusterEventPublisher::new(
516 Self::FALLBACK_CLUSTER_BROADCAST_CAPACITY,
517 ),
518 transcript_publisher: build_transcript_publisher(
522 None,
523 Self::FALLBACK_CLUSTER_BROADCAST_CAPACITY,
524 bounds,
525 ),
526 attempt_owners: crate::worker::AttemptOwnerIndex::new(),
527 cluster_self_node: None,
528 #[cfg(feature = "haematite-backend")]
529 cluster_responder: None,
530 #[cfg(feature = "haematite-backend")]
531 cluster_store: None,
532 #[cfg(feature = "haematite-backend")]
533 watched_peers: Vec::new(),
534 #[cfg(feature = "haematite-backend")]
535 shard_directory: None,
536 #[cfg(feature = "haematite-backend")]
537 request_forwarder: None,
538 #[cfg(feature = "auth")]
539 jwks_cache: None,
540 }),
541 }
542 }
543
544 #[must_use]
546 pub fn namespace_guard(&self) -> &NamespaceGuard {
547 &self.inner.namespace_guard
548 }
549
550 #[must_use]
552 pub fn deploy_guard(&self) -> crate::deploy::DeployGuard {
553 crate::deploy::DeployGuard::new(self.inner.namespace_guard.resolver().clone())
554 }
555
556 #[must_use]
558 pub fn runtime_config(&self) -> &RuntimeConfig {
559 &self.inner.runtime
560 }
561
562 #[must_use]
564 pub fn worker_registry(&self) -> &ConnectedWorkerRegistry {
565 &self.inner.worker_registry
566 }
567
568 #[must_use]
572 pub fn cluster_publisher(&self) -> &crate::cluster_publisher::ClusterEventPublisher {
573 &self.inner.cluster_publisher
574 }
575
576 #[must_use]
581 pub fn transcript_publisher(&self) -> &crate::activity_publisher::ActivityEventPublisher {
582 &self.inner.transcript_publisher
583 }
584
585 #[must_use]
589 pub fn attempt_owners(&self) -> &crate::worker::AttemptOwnerIndex {
590 &self.inner.attempt_owners
591 }
592
593 #[must_use]
605 pub fn intervention_router(&self) -> crate::worker::InterventionRouter {
606 let transport: std::sync::Arc<dyn crate::worker::InterventionTransport> = {
607 #[cfg(feature = "liminal-transport")]
608 {
609 std::sync::Arc::new(crate::worker::LiminalInterventionTransport)
610 }
611 #[cfg(not(feature = "liminal-transport"))]
612 {
613 std::sync::Arc::new(NullInterventionTransport)
614 }
615 };
616 crate::worker::InterventionRouter::new(
617 self.inner.worker_registry.clone(),
618 self.inner.attempt_owners.clone(),
619 transport,
620 )
621 .with_transcript_publisher(self.inner.transcript_publisher.clone())
624 }
625
626 #[must_use]
630 pub fn cluster_self_node(&self) -> Option<&str> {
631 self.inner.cluster_self_node.as_deref()
632 }
633
634 pub fn engine(&self) -> Result<Arc<aion::Engine>, ServerError> {
647 self.inner
648 .namespace_guard
649 .resolver()
650 .engine()
651 .map(Arc::clone)
652 }
653
654 #[must_use]
656 pub fn pending_activities(&self) -> &PendingActivities {
657 &self.inner.pending_activities
658 }
659
660 #[must_use]
662 pub fn heartbeat_tracker(&self) -> &HeartbeatTracker {
663 &self.inner.heartbeat_tracker
664 }
665
666 #[must_use]
668 pub fn drain_state(&self) -> &DrainState {
669 &self.inner.drain_state
670 }
671
672 #[must_use]
674 pub fn metrics(&self) -> Option<&Metrics> {
675 self.inner.metrics.as_ref()
676 }
677
678 #[must_use]
680 pub fn health(&self) -> Option<&HealthState> {
681 self.inner.health.as_ref()
682 }
683
684 #[must_use]
689 pub fn activity_mock_registry(&self) -> Option<&ActivityMockRegistry> {
690 self.inner.activity_mock_registry.as_ref()
691 }
692
693 #[must_use]
699 pub fn outbox_store(&self) -> Option<Arc<dyn OutboxStore>> {
700 self.inner.outbox_store.clone()
701 }
702
703 #[must_use]
712 pub fn namespace_store(&self) -> &Arc<dyn NamespaceStore> {
713 &self.inner.namespace_store
714 }
715
716 #[must_use]
725 pub fn namespace_minter(&self) -> NamespaceMinter {
726 NamespaceMinter::new(
727 Arc::clone(&self.inner.namespace_store),
728 self.inner.runtime.auto_create,
729 )
730 .with_cluster_publisher(self.inner.cluster_publisher.clone())
735 }
736
737 #[must_use]
743 pub fn outbox_wake(&self) -> Arc<tokio::sync::Notify> {
744 Arc::clone(&self.inner.outbox_wake)
745 }
746
747 #[cfg(feature = "haematite-backend")]
753 #[must_use]
754 pub fn is_clustered(&self) -> bool {
755 self.inner.cluster_responder.is_some()
756 }
757
758 #[cfg(feature = "haematite-backend")]
763 #[must_use]
764 pub fn cluster_store(&self) -> Option<&Arc<aion_store_haematite::HaematiteStore>> {
765 self.inner.cluster_store.as_ref()
766 }
767
768 #[cfg(feature = "haematite-backend")]
772 #[must_use]
773 pub fn shard_directory(&self) -> Option<&Arc<crate::routing::StaticShardDirectory>> {
774 self.inner.shard_directory.as_ref()
775 }
776
777 #[cfg(feature = "haematite-backend")]
780 #[must_use]
781 pub fn request_forwarder(&self) -> Option<&Arc<dyn crate::routing::RequestForwarder>> {
782 self.inner.request_forwarder.as_ref()
783 }
784
785 #[must_use]
802 pub fn spawn_heartbeat_sweeper(
803 &self,
804 shutdown: tokio::sync::watch::Receiver<bool>,
805 ) -> tokio::task::JoinHandle<()> {
806 let sweeper = crate::worker::HeartbeatSweeper::new(
807 self.inner.heartbeat_tracker.clone(),
808 self.inner.worker_registry.clone(),
809 self.inner.pending_activities.clone(),
810 self.inner.drain_state.clone(),
811 self.inner.runtime.worker.heartbeat_window,
812 );
813 tokio::spawn(sweeper.run(shutdown))
814 }
815
816 #[cfg(feature = "haematite-backend")]
832 pub fn spawn_cluster_supervisor(
833 &self,
834 config: crate::cluster::SupervisorConfig,
835 shutdown: tokio::sync::watch::Receiver<bool>,
836 ) -> Result<bool, ServerError> {
837 let Some(cluster_store) = self.inner.cluster_store.clone() else {
838 return Ok(false);
839 };
840 if self.inner.watched_peers.is_empty() {
841 return Ok(false);
842 }
843 let engine = Arc::clone(self.inner.namespace_guard.resolver().engine()?);
844 let publisher = Arc::new(self.inner.cluster_publisher.clone());
848 let self_node = self.inner.cluster_self_node.clone().unwrap_or_default();
849 let adopter = Arc::new(crate::cluster::OutboxSettlingAdopter::new(
855 engine,
856 self.inner.outbox_store.clone(),
857 ));
858 let supervisor = crate::cluster::ClusterSupervisor::new(
859 cluster_store,
860 adopter,
861 self.inner.watched_peers.clone(),
862 config,
863 )
864 .with_publisher(publisher, self_node);
865 if !supervisor.watches_any() {
866 return Ok(false);
867 }
868 tokio::spawn(supervisor.run(shutdown));
869 Ok(true)
870 }
871
872 #[cfg(feature = "auth")]
874 #[must_use]
875 pub fn jwks_cache(&self) -> Option<&JwksCache> {
876 self.inner.jwks_cache.as_ref()
877 }
878
879 pub fn shutdown(&self) -> Result<(), ServerError> {
886 self.inner.namespace_guard.resolver().shutdown_engine()
887 }
888}
889
890#[cfg(feature = "auth")]
891async fn build_jwks_cache(runtime: &RuntimeConfig) -> Result<Option<JwksCache>, ServerError> {
892 if !runtime.auth.enabled {
893 return Ok(None);
894 }
895 let Some(url) = runtime.auth.jwks_url.clone() else {
896 return Err(ServerError::Config {
897 message: "auth.jwks_url must not be empty when auth.enabled is true".to_owned(),
898 });
899 };
900 let interval = std::time::Duration::from_secs(runtime.auth.jwks_refresh_seconds);
901 let cache = JwksCache::new(url, interval)
902 .await
903 .map_err(|error| ServerError::Config {
904 message: format!("auth jwks initial fetch failed: {error}"),
905 })?;
906 Ok(Some(cache))
907}
908
909fn metrics_config_error(error: &MetricsError) -> ServerError {
910 ServerError::Config {
911 message: error.to_string(),
912 }
913}
914
915struct EngineAssembly<'a> {
917 instrumented_store: &'a Arc<InstrumentedEventStore>,
919 event_broadcast_capacity: std::num::NonZeroUsize,
921 query_timeout: std::time::Duration,
923 activity_dispatcher: Arc<dyn ActivityDispatcher>,
925 active_registry: Arc<aion::Registry>,
927 bootstrap_coordinator: bool,
929 runtime: &'a RuntimeConfig,
931}
932
933async fn build_engine(assembly: EngineAssembly<'_>) -> Result<aion::Engine, ServerError> {
940 let mut search_attribute_schema = aion_core::SearchAttributeSchema::new();
941 search_attribute_schema
942 .register(
943 crate::namespace::NAMESPACE_ATTRIBUTE,
944 aion_core::SearchAttributeType::String,
945 )
946 .map_err(|error| ServerError::Config {
947 message: format!("failed to register namespace search attribute: {error}"),
948 })?;
949 search_attribute_schema
950 .register(
951 crate::namespace::TASK_QUEUE_ATTRIBUTE,
952 aion_core::SearchAttributeType::String,
953 )
954 .map_err(|error| ServerError::Config {
955 message: format!("failed to register task_queue search attribute: {error}"),
956 })?;
957 let runtime = assembly.runtime;
958 let builder = EngineBuilder::new()
959 .store_arc(assembly.instrumented_store.clone())
960 .event_streaming(assembly.event_broadcast_capacity)
961 .in_memory_visibility()
962 .search_attribute_schema(search_attribute_schema)
963 .scheduler_threads(runtime.scheduler_threads)
964 .outbox_enabled(runtime.outbox.enabled)
965 .activity_dispatcher(assembly.activity_dispatcher)
966 .active_registry(assembly.active_registry)
967 .production_recovery_seam()
968 .signal_router_factory(|runtime: Arc<RuntimeHandle>, handoff| {
969 Arc::new(ConcreteSignalRouter::new(runtime, handoff)) as Arc<dyn SignalRouter>
970 })
971 .query_timeout(assembly.query_timeout)
972 .bootstrap_schedule_coordinator(assembly.bootstrap_coordinator)
978 .load_workflow_sources(runtime.workflow_packages.iter().map(PathBuf::as_path));
979 let builder = if runtime.owned_shards.is_empty() {
985 builder
986 } else {
987 builder.owned_shards(runtime.owned_shards.iter().copied())
988 };
989 builder.build().await.map_err(ServerError::from)
990}
991
992fn required_engine_seams(
997 runtime: &RuntimeConfig,
998) -> Result<(std::num::NonZeroUsize, std::time::Duration), ServerError> {
999 let event_broadcast_capacity = runtime
1000 .websocket
1001 .event_broadcast_capacity
1002 .and_then(std::num::NonZeroUsize::new)
1003 .ok_or_else(|| ServerError::Config {
1004 message: crate::config::EVENT_BROADCAST_CAPACITY_REQUIRED.to_owned(),
1005 })?;
1006 let query_timeout = runtime
1007 .query_timeout
1008 .filter(|timeout| !timeout.is_zero())
1009 .ok_or_else(|| ServerError::Config {
1010 message: crate::config::QUERY_TIMEOUT_REQUIRED.to_owned(),
1011 })?;
1012 Ok((event_broadcast_capacity, query_timeout))
1013}
1014
1015fn install_outbox_delivery(
1022 pending_activities: &PendingActivities,
1023 engine: &Arc<aion::Engine>,
1024 outbox_enabled: bool,
1025) {
1026 if outbox_enabled {
1027 let callback = Arc::new(crate::worker::ServerOutboxDeliveryCallback::new(
1028 Arc::clone(engine),
1029 ));
1030 pending_activities.set_outbox_delivery(callback);
1031 }
1032}
1033
1034fn build_bridge_dispatcher(
1047 runtime: &RuntimeConfig,
1048 worker_registry: &ConnectedWorkerRegistry,
1049 pending_activities: &PendingActivities,
1050 heartbeat_tracker: &HeartbeatTracker,
1051 drain_state: &DrainState,
1052) -> (WorkerActivityDispatcher, crate::worker::AttemptOwnerIndex) {
1053 let attempt_owners = crate::worker::AttemptOwnerIndex::new();
1054 let dispatcher = WorkerActivityDispatcher::new(
1055 worker_registry.clone(),
1056 runtime.default_namespace.clone(),
1057 heartbeat_tracker.clone(),
1058 )
1059 .with_pending(pending_activities.clone())
1060 .with_drain_state(drain_state.clone())
1061 .with_tokio_handle(tokio::runtime::Handle::current())
1062 .with_attempt_owners(attempt_owners.clone());
1063 (dispatcher, attempt_owners)
1064}
1065
1066fn decorate_activity_dispatcher(
1067 dispatcher: WorkerActivityDispatcher,
1068 dev_enabled: bool,
1069) -> (Arc<dyn ActivityDispatcher>, Option<ActivityMockRegistry>) {
1070 if dev_enabled {
1071 let registry = ActivityMockRegistry::new();
1072 let decorated = DevMockingDispatcher::new(Arc::new(dispatcher), registry.clone());
1073 (Arc::new(decorated), Some(registry))
1074 } else {
1075 (Arc::new(dispatcher), None)
1076 }
1077}
1078
1079fn required_cluster_broadcast_capacity(
1084 runtime: &RuntimeConfig,
1085) -> Result<std::num::NonZeroUsize, ServerError> {
1086 runtime
1087 .websocket
1088 .cluster_broadcast_capacity
1089 .and_then(std::num::NonZeroUsize::new)
1090 .ok_or_else(|| ServerError::Config {
1091 message: crate::config::CLUSTER_BROADCAST_CAPACITY_REQUIRED.to_owned(),
1092 })
1093}
1094
1095fn build_real_time_publishers(
1108 runtime: &RuntimeConfig,
1109 observability_store: Option<Arc<dyn aion_store::ObservabilityStore>>,
1110) -> Result<
1111 (
1112 crate::cluster_publisher::ClusterEventPublisher,
1113 crate::activity_publisher::ActivityEventPublisher,
1114 ),
1115 ServerError,
1116> {
1117 let capacity = required_cluster_broadcast_capacity(runtime)?;
1118 Ok((
1119 crate::cluster_publisher::ClusterEventPublisher::new(capacity),
1120 build_transcript_publisher(observability_store, capacity, transcript_bounds(runtime)),
1121 ))
1122}
1123
1124fn transcript_bounds(runtime: &RuntimeConfig) -> crate::activity_bounds::TranscriptBounds {
1126 crate::activity_bounds::TranscriptBounds {
1127 max_event_bytes: runtime.observability.max_event_bytes,
1128 max_stream_events: runtime.observability.max_stream_events,
1129 }
1130}
1131
1132fn build_transcript_publisher(
1143 observability_store: Option<Arc<dyn aion_store::ObservabilityStore>>,
1144 capacity: std::num::NonZeroUsize,
1145 bounds: crate::activity_bounds::TranscriptBounds,
1146) -> crate::activity_publisher::ActivityEventPublisher {
1147 let store = observability_store
1148 .unwrap_or_else(|| Arc::new(aion_store::InMemoryObservabilityStore::default()));
1149 crate::activity_publisher::ActivityEventPublisher::new(store, capacity).with_bounds(bounds)
1150}
1151
1152#[cfg(feature = "haematite-backend")]
1154struct RoutingState {
1155 shard_directory: Option<Arc<crate::routing::StaticShardDirectory>>,
1156 request_forwarder: Option<Arc<dyn crate::routing::RequestForwarder>>,
1157}
1158
1159#[cfg(feature = "haematite-backend")]
1163fn build_routing_state(
1164 cluster_store: Option<&Arc<aion_store_haematite::HaematiteStore>>,
1165 directory_peers: Vec<crate::routing::DirectoryPeer>,
1166 self_node_id: Option<String>,
1167) -> RoutingState {
1168 let Some(store) = cluster_store else {
1169 return RoutingState {
1170 shard_directory: None,
1171 request_forwarder: None,
1172 };
1173 };
1174 RoutingState {
1175 shard_directory: Some(Arc::new(crate::routing::StaticShardDirectory::new(
1176 Arc::clone(store),
1177 directory_peers,
1178 self_node_id,
1179 ))),
1180 request_forwarder: Some(Arc::new(crate::routing::GrpcRequestForwarder::new())),
1181 }
1182}
1183
1184struct ConnectedStore {
1194 event_store: Arc<dyn EventStore>,
1195 outbox_store: Option<Arc<dyn OutboxStore>>,
1196 namespace_store: Arc<dyn NamespaceStore>,
1203 observability_store: Option<Arc<dyn aion_store::ObservabilityStore>>,
1211 bootstrap_coordinator: bool,
1212 #[cfg(feature = "haematite-backend")]
1213 cluster_responder: Option<aion_store_haematite::ClusterResponder>,
1214 #[cfg(feature = "haematite-backend")]
1218 cluster_store: Option<Arc<aion_store_haematite::HaematiteStore>>,
1219 #[cfg(feature = "haematite-backend")]
1222 watched_peers: Vec<crate::cluster::WatchedPeer>,
1223 #[cfg(feature = "haematite-backend")]
1227 directory_peers: Vec<crate::routing::DirectoryPeer>,
1228 #[cfg(feature = "haematite-backend")]
1232 self_node_id: Option<String>,
1233}
1234
1235impl ConnectedStore {
1236 fn local(
1243 event_store: Arc<dyn EventStore>,
1244 outbox_store: Option<Arc<dyn OutboxStore>>,
1245 namespace_store: Arc<dyn NamespaceStore>,
1246 ) -> Self {
1247 Self {
1248 event_store,
1249 outbox_store,
1250 namespace_store,
1251 observability_store: None,
1255 bootstrap_coordinator: true,
1256 #[cfg(feature = "haematite-backend")]
1257 cluster_responder: None,
1258 #[cfg(feature = "haematite-backend")]
1259 cluster_store: None,
1260 #[cfg(feature = "haematite-backend")]
1261 watched_peers: Vec::new(),
1262 #[cfg(feature = "haematite-backend")]
1263 directory_peers: Vec::new(),
1264 #[cfg(feature = "haematite-backend")]
1265 self_node_id: None,
1266 }
1267 }
1268}
1269
1270async fn connect_store(config: StoreConfig) -> Result<ConnectedStore, ServerError> {
1280 match config.backend {
1281 StoreBackend::Memory => {
1282 let leaf = Arc::new(aion_store::InMemoryStore::default());
1285 let namespace_store: Arc<dyn NamespaceStore> = leaf.clone();
1286 Ok(ConnectedStore::local(leaf, None, namespace_store))
1287 }
1288 StoreBackend::LibSql => {
1289 #[cfg(feature = "libsql-backend")]
1290 {
1291 connect_libsql_store(config).await
1292 }
1293 #[cfg(not(feature = "libsql-backend"))]
1294 {
1295 let _ = config;
1296 connect_libsql_store_unavailable()
1297 }
1298 }
1299 StoreBackend::Haematite => {
1300 #[cfg(feature = "haematite-backend")]
1301 {
1302 connect_haematite_store(config).await
1303 }
1304 #[cfg(not(feature = "haematite-backend"))]
1305 {
1306 let _ = config;
1307 connect_haematite_store_unavailable()
1308 }
1309 }
1310 }
1311}
1312
1313#[cfg(feature = "libsql-backend")]
1317async fn connect_libsql_store(config: StoreConfig) -> Result<ConnectedStore, ServerError> {
1318 let Some(url) = config.url else {
1319 return Err(ServerError::Config {
1320 message: "store.url must not be empty when store.backend is libsql".to_owned(),
1321 });
1322 };
1323 let store = LibSqlStore::open(url.clone())
1324 .await
1325 .map_err(ServerError::from)?;
1326 store
1327 .validate_event_compatibility()
1328 .await
1329 .map_err(|error| match error {
1330 aion_store::StoreError::Serialization(_) => ServerError::Config {
1331 message: format!(
1332 "Database schema mismatch — delete {url} and restart, or run migrations."
1333 ),
1334 },
1335 other => ServerError::from(other),
1336 })?;
1337 let leaf = Arc::new(store);
1338 let event_store: Arc<dyn EventStore> = leaf.clone();
1339 let namespace_store: Arc<dyn NamespaceStore> = leaf.clone();
1340 let outbox_store: Arc<dyn OutboxStore> = leaf;
1341 Ok(ConnectedStore::local(
1342 event_store,
1343 Some(outbox_store),
1344 namespace_store,
1345 ))
1346}
1347
1348#[cfg(not(feature = "libsql-backend"))]
1352fn connect_libsql_store_unavailable() -> Result<ConnectedStore, ServerError> {
1353 Err(ServerError::Config {
1354 message: "store.backend = libsql requires the aion-server `libsql-backend` feature"
1355 .to_owned(),
1356 })
1357}
1358
1359#[cfg(feature = "haematite-backend")]
1379async fn connect_haematite_store(config: StoreConfig) -> Result<ConnectedStore, ServerError> {
1380 let Some(data_dir) = config.data_dir else {
1381 return Err(ServerError::Config {
1382 message: "store.data_dir must not be empty when store.backend is haematite".to_owned(),
1383 });
1384 };
1385 let shard_count = config.shard_count;
1386 let owned_shards = config.owned_shards.clone();
1387 let cluster = config.cluster.clone();
1388 let watched_peers: Vec<crate::cluster::WatchedPeer> = cluster
1393 .as_ref()
1394 .map(|cluster| {
1395 cluster
1396 .peers
1397 .iter()
1398 .map(|peer| crate::cluster::WatchedPeer {
1399 name: peer.name.clone(),
1400 owned_shards: peer.owned_shards.clone(),
1401 })
1402 .collect()
1403 })
1404 .unwrap_or_default();
1405 let directory_peers: Vec<crate::routing::DirectoryPeer> = cluster
1408 .as_ref()
1409 .map(|cluster| {
1410 cluster
1411 .peers
1412 .iter()
1413 .map(|peer| crate::routing::DirectoryPeer {
1414 name: peer.name.clone(),
1415 owned_shards: peer.owned_shards.clone(),
1416 grpc_addr: peer.grpc_address,
1417 })
1418 .collect()
1419 })
1420 .unwrap_or_default();
1421 let self_node_id: Option<String> = cluster.as_ref().map(|cluster| cluster.node_id.clone());
1424 let (store, responder) =
1428 tokio::task::spawn_blocking(move || build_haematite_store(&data_dir, shard_count, cluster))
1429 .await
1430 .map_err(|error| ServerError::Config {
1431 message: format!("haematite store initialization task failed: {error}"),
1432 })??;
1433
1434 let bootstrap_coordinator = if owned_shards.is_empty() {
1438 true
1439 } else {
1440 store.set_owned_shards(owned_shards.iter().copied());
1441 store.owns_workflow_shard(&aion::schedule_coordinator_workflow_id())
1442 };
1443
1444 let leaf = Arc::new(store);
1445 let event_store: Arc<dyn EventStore> = leaf.clone();
1446 let outbox_store: Arc<dyn OutboxStore> = leaf.clone();
1447 let namespace_store: Arc<dyn NamespaceStore> = leaf.clone();
1451 let observability_store: Arc<dyn aion_store::ObservabilityStore> = leaf.clone();
1456 let cluster_store = responder.as_ref().map(|_| leaf);
1460 let (watched_peers, directory_peers, self_node_id) = if cluster_store.is_some() {
1461 (watched_peers, directory_peers, self_node_id)
1462 } else {
1463 (Vec::new(), Vec::new(), None)
1464 };
1465 Ok(ConnectedStore {
1466 event_store,
1467 outbox_store: Some(outbox_store),
1468 namespace_store,
1469 observability_store: Some(observability_store),
1470 bootstrap_coordinator,
1471 cluster_responder: responder,
1472 cluster_store,
1473 watched_peers,
1474 directory_peers,
1475 self_node_id,
1476 })
1477}
1478
1479#[cfg(feature = "haematite-backend")]
1494fn build_haematite_store(
1495 data_dir: &str,
1496 shard_count: usize,
1497 cluster: Option<crate::config::ClusterConfig>,
1498) -> Result<
1499 (
1500 aion_store_haematite::HaematiteStore,
1501 Option<aion_store_haematite::ClusterResponder>,
1502 ),
1503 ServerError,
1504> {
1505 build_haematite_store_with_hook(data_dir, shard_count, cluster, || Ok(()))
1506}
1507
1508#[cfg(feature = "haematite-backend")]
1509fn build_haematite_store_with_hook(
1510 data_dir: &str,
1511 shard_count: usize,
1512 cluster: Option<crate::config::ClusterConfig>,
1513 before_backend_touch: impl FnOnce() -> Result<(), std::io::Error>,
1514) -> Result<
1515 (
1516 aion_store_haematite::HaematiteStore,
1517 Option<aion_store_haematite::ClusterResponder>,
1518 ),
1519 ServerError,
1520> {
1521 use aion_store_haematite::{ClusterBootstrap, HaematiteStore};
1522
1523 let private_root = crate::filesystem::ConfinedDir::open_or_create(std::path::Path::new(
1527 data_dir,
1528 ))
1529 .map_err(|error| ServerError::Config {
1530 message: format!("unsafe store.data_dir `{data_dir}`: {error}"),
1531 })?;
1532
1533 for shard in 0..shard_count {
1537 private_root
1538 .create_dir_all(std::path::Path::new(&format!("shard-{shard}")))
1539 .map_err(|error| ServerError::Config {
1540 message: format!(
1541 "failed to materialize shard-{shard} under store.data_dir `{data_dir}`: {error}"
1542 ),
1543 })?;
1544 }
1545 private_root
1546 .harden_tree()
1547 .map_err(|error| private_store_mode_error(data_dir, &error))?;
1548
1549 before_backend_touch().map_err(|error| ServerError::Config {
1552 message: format!("store.data_dir pre-open hook failed: {error}"),
1553 })?;
1554
1555 #[cfg(unix)]
1556 let backend_path = private_root
1557 .backend_path()
1558 .map_err(|error| ServerError::Config {
1559 message: format!("failed to resolve held store.data_dir `{data_dir}`: {error}"),
1560 })?;
1561 #[cfg(all(unix, not(any(target_os = "linux", target_os = "android"))))]
1562 crate::filesystem::validate_ambient_backend_ancestors(&backend_path).map_err(|error| {
1563 let (component, reason) = error.into_parts();
1564 ServerError::UnsafeDataRootAncestor {
1565 data_root: backend_path.clone(),
1566 component,
1567 reason,
1568 }
1569 })?;
1570 #[cfg(not(unix))]
1571 let backend_path = std::path::PathBuf::from(data_dir);
1572
1573 let Some(cluster) = cluster else {
1574 let store = if backend_path.join("config.json").exists() {
1575 HaematiteStore::open(&backend_path).map_err(ServerError::from)?
1576 } else {
1577 HaematiteStore::create_with_shard_count(&backend_path, shard_count)
1578 .map_err(ServerError::from)?
1579 };
1580 store.materialize_all_shards().map_err(ServerError::from)?;
1581 private_root
1582 .harden_tree()
1583 .map_err(|error| private_store_mode_error(data_dir, &error))?;
1584 let store = store.retain_data_root_capability(private_root);
1585 return Ok((store, None));
1586 };
1587
1588 let boot = ClusterBootstrap {
1589 node_id: cluster.node_id,
1590 bind_address: cluster.bind_address,
1591 members: cluster.members,
1592 peers: cluster
1593 .peers
1594 .into_iter()
1595 .map(|peer| (peer.name, peer.address))
1596 .collect(),
1597 timeout: HAEMATITE_CLUSTER_OP_TIMEOUT,
1598 };
1599 let (store, responder) =
1600 HaematiteStore::open_or_create_distributed(&backend_path, shard_count, boot)
1601 .map_err(ServerError::from)?;
1602 store.materialize_all_shards().map_err(ServerError::from)?;
1603 private_root
1604 .harden_tree()
1605 .map_err(|error| private_store_mode_error(data_dir, &error))?;
1606 let store = store.retain_data_root_capability(private_root);
1607 Ok((store, Some(responder)))
1608}
1609
1610#[cfg(feature = "haematite-backend")]
1611fn private_store_mode_error(data_dir: &str, error: &std::io::Error) -> ServerError {
1612 ServerError::Config {
1613 message: format!(
1614 "failed to apply private modes under store.data_dir `{data_dir}`: {error}"
1615 ),
1616 }
1617}
1618
1619#[cfg(feature = "haematite-backend")]
1621const HAEMATITE_CLUSTER_OP_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(5);
1622
1623#[cfg(not(feature = "haematite-backend"))]
1627fn connect_haematite_store_unavailable() -> Result<ConnectedStore, ServerError> {
1628 Err(ServerError::Config {
1629 message: "store.backend = haematite requires the aion-server `haematite-backend` feature"
1630 .to_owned(),
1631 })
1632}
1633
1634#[cfg(not(feature = "liminal-transport"))]
1643#[derive(Clone, Debug)]
1644struct NullInterventionTransport;
1645
1646#[cfg(not(feature = "liminal-transport"))]
1647#[async_trait::async_trait]
1648impl crate::worker::InterventionTransport for NullInterventionTransport {
1649 async fn push(
1650 &self,
1651 _worker: &crate::worker::WorkerHandle,
1652 _command: aion_core::InterventionCommand,
1653 ) -> Result<aion_core::InterventionOutcome, ServerError> {
1654 Err(ServerError::worker_connection_lost(
1655 "intervention",
1656 "no intervention push transport is compiled in".to_owned(),
1657 ))
1658 }
1659}
1660
1661#[cfg(test)]
1662mod tests {
1663 use std::{net::SocketAddr, time::Duration};
1664
1665 use aion_store::InMemoryStore;
1666
1667 use super::ServerState;
1668 use crate::config::{
1669 AuthConfig, AuthoringConfig, DeployConfig, DevConfig, ListenConfig, MetricsConfig,
1670 NamespaceConfig, NamespaceMode, OpsConsoleAssetSource, OpsConsoleConfig, OutboxConfig,
1671 RuntimeConfig, WebSocketConfig, WorkerConfig,
1672 };
1673
1674 fn runtime_config() -> RuntimeConfig {
1675 RuntimeConfig {
1676 listen: ListenConfig {
1677 grpc: SocketAddr::from(([127, 0, 0, 1], 50051)),
1678 http: SocketAddr::from(([127, 0, 0, 1], 8080)),
1679 },
1680 tls: None,
1681 auth: AuthConfig {
1682 enabled: false,
1683 jwks_url: None,
1684 jwks_refresh_seconds: 300,
1685 },
1686 ops_console: OpsConsoleConfig {
1687 source: OpsConsoleAssetSource::Embedded,
1688 },
1689 namespace: NamespaceConfig {
1690 mode: NamespaceMode::SharedEngine,
1691 },
1692 worker: WorkerConfig {
1693 heartbeat_window: Duration::from_millis(30_000),
1694 },
1695 websocket: WebSocketConfig {
1696 outbound_buffer_bound: 32,
1697 event_broadcast_capacity: Some(64),
1698 cluster_broadcast_capacity: Some(64),
1699 },
1700 workflow_packages: Vec::new(),
1701 deploy: DeployConfig::default(),
1702 authoring: AuthoringConfig::default(),
1703 dev: DevConfig::default(),
1704 outbox: OutboxConfig::default(),
1705 observability: crate::config::ObservabilityConfig::default(),
1706 scheduler_threads: 1,
1707 query_timeout: Some(Duration::from_millis(10_000)),
1708 default_namespace: "default".to_owned(),
1709 auto_create: crate::config::AutoCreate::Open,
1710 max_in_flight_activities: crate::config::DEFAULT_MAX_IN_FLIGHT_ACTIVITIES,
1711 drain_timeout: Duration::from_secs(30),
1712 metrics: MetricsConfig { enabled: true },
1713 owned_shards: Vec::new(),
1714 cors_allowed_origins: Vec::new(),
1715 }
1716 }
1717
1718 #[tokio::test]
1719 async fn builds_state_with_in_memory_store() -> Result<(), Box<dyn std::error::Error>> {
1720 let state =
1721 ServerState::build_with_store(InMemoryStore::default(), runtime_config()).await?;
1722
1723 std::hint::black_box(state.namespace_guard());
1724 std::hint::black_box(state.worker_registry());
1725
1726 Ok(())
1727 }
1728
1729 #[tokio::test]
1730 async fn namespace_store_is_reachable_and_functional_after_default_boot()
1731 -> Result<(), Box<dyn std::error::Error>> {
1732 use aion_store::{MintOutcome, NamespaceOrigin};
1733
1734 let state =
1738 ServerState::build_with_store(InMemoryStore::default(), runtime_config()).await?;
1739
1740 let store = state.namespace_store();
1741
1742 let outcome = store
1744 .register_namespace("orders", NamespaceOrigin::WorkerMint)
1745 .await?;
1746 assert_eq!(
1747 outcome,
1748 MintOutcome::Created,
1749 "the first reference to a namespace mints it"
1750 );
1751
1752 let again = store
1754 .register_namespace("orders", NamespaceOrigin::WorkerMint)
1755 .await?;
1756 assert_eq!(
1757 again,
1758 MintOutcome::AlreadyExisted,
1759 "a second reference touches the existing record rather than re-creating it"
1760 );
1761
1762 let fetched = store.get_namespace("orders").await?;
1764 let record = fetched.ok_or("registered namespace must be retrievable via get_namespace")?;
1765 assert_eq!(record.name, "orders");
1766 assert_eq!(record.origin, NamespaceOrigin::WorkerMint);
1767
1768 let listed = store.list_namespaces().await?;
1770 assert!(
1771 listed.iter().any(|record| record.name == "orders"),
1772 "list_namespaces returns the minted namespace"
1773 );
1774
1775 Ok(())
1776 }
1777
1778 #[cfg(feature = "haematite-backend")]
1779 #[tokio::test(flavor = "multi_thread")]
1780 async fn connect_store_haematite_round_trips_through_event_store()
1781 -> Result<(), Box<dyn std::error::Error>> {
1782 use aion_core::{ContentType, EventEnvelope, PackageVersion, Payload, RunId, WorkflowId};
1783 use aion_store::WriteToken;
1784 use chrono::Utc;
1785
1786 use crate::config::{StoreBackend, StoreConfig};
1787
1788 let data_dir = crate::test_support::private_tempdir()?;
1789 let connected = super::connect_store(StoreConfig {
1793 backend: StoreBackend::Haematite,
1794 url: None,
1795 owned_shards: Vec::new(),
1796 data_dir: Some(data_dir.path().to_string_lossy().into_owned()),
1797 shard_count: 1,
1798 cluster: None,
1799 })
1800 .await?;
1801 let event_store = connected.event_store;
1802 assert!(
1803 connected.outbox_store.is_some(),
1804 "the haematite backend shares its leaf store as the dispatcher's outbox store"
1805 );
1806 assert!(
1807 connected.bootstrap_coordinator,
1808 "a single-node haematite boot owns all shards and bootstraps the coordinator"
1809 );
1810 assert!(
1811 connected.cluster_responder.is_none(),
1812 "a single-node (no [cluster]) haematite boot has no distributed responder"
1813 );
1814
1815 let workflow_id = WorkflowId::new_v4();
1816 let event = aion_core::Event::WorkflowStarted {
1817 envelope: EventEnvelope {
1818 seq: 1,
1819 recorded_at: Utc::now(),
1820 workflow_id: workflow_id.clone(),
1821 },
1822 workflow_type: String::from("checkout"),
1823 input: Payload::new(ContentType::Json, b"{}".to_vec()),
1824 run_id: RunId::new_v4(),
1825 parent_run_id: None,
1826 package_version: PackageVersion::new("a".repeat(64)),
1827 };
1828 event_store
1829 .append(
1830 WriteToken::recorder(),
1831 &workflow_id,
1832 std::slice::from_ref(&event),
1833 0,
1834 )
1835 .await?;
1836 let history = event_store.read_history(&workflow_id).await?;
1837 assert_eq!(
1838 history.len(),
1839 1,
1840 "an event appended through the server's dyn EventStore reads back"
1841 );
1842 Ok(())
1843 }
1844
1845 #[cfg(all(feature = "haematite-backend", unix))]
1846 #[test]
1847 fn haematite_root_swap_before_first_backend_touch_cannot_redirect_writes()
1848 -> Result<(), Box<dyn std::error::Error>> {
1849 use std::os::unix::fs::symlink;
1850
1851 let sandbox = crate::test_support::private_tempdir()?;
1852 let configured_root = sandbox.path().join("data");
1853 let held_root = sandbox.path().join("held-data");
1854 let outside = sandbox.path().join("outside");
1855 std::fs::create_dir(&outside)?;
1856 let configured = configured_root
1857 .to_str()
1858 .ok_or("temporary data path was not UTF-8")?;
1859
1860 let (store, responder) =
1861 super::build_haematite_store_with_hook(configured, 4, None, || {
1862 std::fs::rename(&configured_root, &held_root)?;
1867 symlink(&outside, &configured_root)?;
1868 Ok(())
1869 })?;
1870 assert!(responder.is_none());
1871
1872 let outside_entries = std::fs::read_dir(&outside)?.collect::<Result<Vec<_>, _>>()?;
1873 assert!(
1874 outside_entries.is_empty(),
1875 "Haematite followed the replaced ambient root and wrote outside"
1876 );
1877 assert!(held_root.join("config.json").is_file());
1878 for shard in 0..4 {
1879 let shard_path = held_root.join(format!("shard-{shard}"));
1880 assert!(shard_path.is_dir(), "shard {shard} was not materialized");
1881 assert!(
1882 std::fs::read_dir(&shard_path)?
1883 .next()
1884 .transpose()?
1885 .is_some(),
1886 "shard {shard} did not run Haematite's materialization path"
1887 );
1888 }
1889
1890 drop(store);
1891 Ok(())
1892 }
1893
1894 #[cfg(all(
1895 feature = "haematite-backend",
1896 any(target_os = "linux", target_os = "android")
1897 ))]
1898 #[tokio::test]
1899 async fn proc_fd_backend_path_survives_a_post_startup_root_swap()
1900 -> Result<(), Box<dyn std::error::Error>> {
1901 use std::os::unix::fs::symlink;
1902
1903 use aion_core::{ContentType, EventEnvelope, PackageVersion, Payload, RunId, WorkflowId};
1904 use aion_store::{WritableEventStore as _, WriteToken};
1905 use chrono::Utc;
1906
1907 let sandbox = crate::test_support::private_tempdir()?;
1908 let configured_root = sandbox.path().join("data");
1909 let held_root = sandbox.path().join("held-data");
1910 let capture = sandbox.path().join("capture");
1911 std::fs::create_dir(&capture)?;
1912 let configured = configured_root
1913 .to_str()
1914 .ok_or("temporary data path was not UTF-8")?;
1915
1916 let (store, responder) = super::build_haematite_store(configured, 4, None)?;
1917 assert!(responder.is_none());
1918 std::fs::rename(&configured_root, &held_root)?;
1919 symlink(&capture, &configured_root)?;
1920
1921 let workflow_id = WorkflowId::new_v4();
1922 let event = aion_core::Event::WorkflowStarted {
1923 envelope: EventEnvelope {
1924 seq: 1,
1925 recorded_at: Utc::now(),
1926 workflow_id: workflow_id.clone(),
1927 },
1928 workflow_type: String::from("post-startup-root-swap"),
1929 input: Payload::new(ContentType::Json, b"{}".to_vec()),
1930 run_id: RunId::new_v4(),
1931 parent_run_id: None,
1932 package_version: PackageVersion::new("a".repeat(64)),
1933 };
1934 store
1935 .append(
1936 WriteToken::recorder(),
1937 &workflow_id,
1938 std::slice::from_ref(&event),
1939 0,
1940 )
1941 .await?;
1942
1943 let captured = std::fs::read_dir(&capture)?.collect::<Result<Vec<_>, _>>()?;
1944 assert!(
1945 captured.is_empty(),
1946 "post-startup append followed the replacement symlink into capture"
1947 );
1948 assert!(held_root.join("config.json").is_file());
1949 drop(store);
1950 Ok(())
1951 }
1952
1953 #[cfg(all(
1954 feature = "haematite-backend",
1955 unix,
1956 not(any(target_os = "linux", target_os = "android"))
1957 ))]
1958 #[test]
1959 fn path_ambient_haematite_refuses_group_or_world_writable_ancestors()
1960 -> Result<(), Box<dyn std::error::Error>> {
1961 use std::os::unix::fs::PermissionsExt as _;
1962
1963 let sandbox = crate::test_support::private_tempdir()?;
1964 std::fs::set_permissions(sandbox.path(), std::fs::Permissions::from_mode(0o700))?;
1965
1966 for mode in [0o770, 0o1777] {
1967 let shared = sandbox.path().join(format!("shared-{mode:o}"));
1968 let data_root = shared.join("data");
1969 std::fs::create_dir(&shared)?;
1970 std::fs::set_permissions(&shared, std::fs::Permissions::from_mode(mode))?;
1971 std::fs::create_dir(&data_root)?;
1972 std::fs::set_permissions(&data_root, std::fs::Permissions::from_mode(0o700))?;
1973 let configured = data_root
1974 .to_str()
1975 .ok_or("temporary data path was not UTF-8")?;
1976
1977 let Err(error) = super::build_haematite_store(configured, 4, None) else {
1978 return Err(format!("mode {mode:04o} ancestor was accepted").into());
1979 };
1980 let message = error.to_string();
1981 let crate::ServerError::UnsafeDataRootAncestor {
1982 data_root: resolved_root,
1983 component,
1984 reason,
1985 } = error
1986 else {
1987 return Err(format!("expected typed unsafe-ancestor error, got {message}").into());
1988 };
1989 assert_eq!(resolved_root, std::fs::canonicalize(&data_root)?);
1990 assert_eq!(component, std::fs::canonicalize(&shared)?);
1991 assert!(
1992 reason.contains(&format!("mode {mode:04o}")),
1993 "unexpected reason: {reason}"
1994 );
1995 if mode & 0o1000 != 0 {
1996 assert!(reason.contains("sticky bit is not accepted"));
1997 }
1998 assert!(message.contains("private Aion home"));
1999 assert!(
2000 !data_root.join("config.json").exists(),
2001 "Haematite touched its ambient path before the refusal"
2002 );
2003 }
2004 Ok(())
2005 }
2006
2007 #[cfg(all(feature = "haematite-backend", target_os = "macos"))]
2008 #[test]
2009 fn path_ambient_haematite_refuses_mutating_allow_acl_ancestor()
2010 -> Result<(), Box<dyn std::error::Error>> {
2011 use std::os::unix::fs::PermissionsExt as _;
2012
2013 let sandbox = crate::test_support::private_tempdir()?;
2014 std::fs::set_permissions(sandbox.path(), std::fs::Permissions::from_mode(0o700))?;
2015 let shared = sandbox.path().join("acl-shared");
2016 let data_root = shared.join("data");
2017 std::fs::create_dir(&shared)?;
2018 std::fs::set_permissions(&shared, std::fs::Permissions::from_mode(0o700))?;
2019 let acl = "everyone allow list,search,add_file,add_subdirectory,delete_child";
2020 let status = std::process::Command::new("chmod")
2021 .arg("+a")
2022 .arg(acl)
2023 .arg(&shared)
2024 .status()?;
2025 assert!(status.success(), "failed to install Darwin regression ACL");
2026 let configured = data_root
2027 .to_str()
2028 .ok_or("temporary data path was not UTF-8")?;
2029
2030 let result = super::build_haematite_store(configured, 4, None);
2031 let cleanup = std::process::Command::new("chmod")
2032 .arg("-RN")
2033 .arg(&shared)
2034 .status()?;
2035 assert!(cleanup.success(), "failed to clean Darwin regression ACL");
2036
2037 let Err(error) = result else {
2038 return Err("mutating non-euid allow ACL ancestor was accepted".into());
2039 };
2040 let message = error.to_string();
2041 let crate::ServerError::UnsafeDataRootAncestor {
2042 component, reason, ..
2043 } = error
2044 else {
2045 return Err(format!("expected typed unsafe-ancestor error, got {message}").into());
2046 };
2047 assert_eq!(component, std::fs::canonicalize(&shared)?);
2048 assert!(
2049 reason.contains("allow"),
2050 "reason did not name the ACE: {reason}"
2051 );
2052 assert!(
2053 reason.contains("everyone"),
2054 "reason did not name the ACE principal: {reason}"
2055 );
2056 assert!(
2057 !data_root.join("config.json").exists(),
2058 "Haematite touched its ambient path before the ACL refusal"
2059 );
2060 Ok(())
2061 }
2062
2063 #[cfg(all(feature = "haematite-backend", target_os = "macos"))]
2064 #[test]
2065 fn path_ambient_haematite_accepts_the_euid_uuid_allow_ace()
2066 -> Result<(), Box<dyn std::error::Error>> {
2067 use std::os::unix::fs::PermissionsExt as _;
2068
2069 use exacl::{AclEntry, AclOption, Perm};
2070
2071 let sandbox = crate::test_support::private_tempdir()?;
2072 std::fs::set_permissions(sandbox.path(), std::fs::Permissions::from_mode(0o700))?;
2073 let private_parent = sandbox.path().join("euid-uuid-allow");
2074 let data_root = private_parent.join("data");
2075 std::fs::create_dir(&private_parent)?;
2076 std::fs::set_permissions(&private_parent, std::fs::Permissions::from_mode(0o700))?;
2077
2078 let server_uid = rustix::process::geteuid().as_raw();
2079 let ace_qualifier = crate::filesystem::darwin_user_uuid_for_test(server_uid)?;
2080 let entry = AclEntry::allow_user(
2081 &ace_qualifier.to_string(),
2082 Perm::EXECUTE | Perm::WRITE | Perm::APPEND | Perm::DELETE_CHILD,
2083 None,
2084 );
2085 exacl::setfacl(
2086 &[private_parent.as_path()],
2087 &[entry],
2088 AclOption::SYMLINK_ACL,
2089 )?;
2090 let configured = data_root
2091 .to_str()
2092 .ok_or("temporary data path was not UTF-8")?;
2093
2094 let result = super::build_haematite_store(configured, 4, None);
2095 let cleanup = std::process::Command::new("chmod")
2096 .arg("-RN")
2097 .arg(&private_parent)
2098 .status()?;
2099 assert!(cleanup.success(), "failed to clean euid UUID allow ACL");
2100
2101 let (store, responder) = result?;
2102 assert!(responder.is_none());
2103 assert!(data_root.join("config.json").is_file());
2104 drop(store);
2105 Ok(())
2106 }
2107
2108 #[cfg(all(feature = "haematite-backend", target_os = "macos"))]
2109 #[test]
2110 fn path_ambient_haematite_refuses_a_non_euid_user_uuid_allow_ace()
2111 -> Result<(), Box<dyn std::error::Error>> {
2112 use std::os::unix::fs::PermissionsExt as _;
2113
2114 use exacl::{AclEntry, AclOption, Perm};
2115
2116 let sandbox = crate::test_support::private_tempdir()?;
2117 std::fs::set_permissions(sandbox.path(), std::fs::Permissions::from_mode(0o700))?;
2118 let shared = sandbox.path().join("non-euid-uuid-allow");
2119 let data_root = shared.join("data");
2120 std::fs::create_dir(&shared)?;
2121 std::fs::set_permissions(&shared, std::fs::Permissions::from_mode(0o700))?;
2122
2123 let server_uid = rustix::process::geteuid().as_raw();
2124 let foreign_uid = u32::from(server_uid == 0);
2125 let foreign_qualifier = crate::filesystem::darwin_user_uuid_for_test(foreign_uid)?;
2126 let entry = AclEntry::allow_user(
2127 &foreign_qualifier.to_string(),
2128 Perm::EXECUTE | Perm::WRITE | Perm::APPEND | Perm::DELETE_CHILD,
2129 None,
2130 );
2131 exacl::setfacl(&[shared.as_path()], &[entry], AclOption::SYMLINK_ACL)?;
2132 let configured = data_root
2133 .to_str()
2134 .ok_or("temporary data path was not UTF-8")?;
2135
2136 let result = super::build_haematite_store(configured, 4, None);
2137 let cleanup = std::process::Command::new("chmod")
2138 .arg("-RN")
2139 .arg(&shared)
2140 .status()?;
2141 assert!(cleanup.success(), "failed to clean non-euid UUID allow ACL");
2142
2143 let Err(error) = result else {
2144 return Err("mutating non-euid user UUID allow ACE was accepted".into());
2145 };
2146 let message = error.to_string();
2147 let crate::ServerError::UnsafeDataRootAncestor {
2148 component, reason, ..
2149 } = error
2150 else {
2151 return Err(format!("expected typed unsafe-ancestor error, got {message}").into());
2152 };
2153 assert_eq!(component, std::fs::canonicalize(&shared)?);
2154 assert!(
2155 reason.contains("allow") && reason.contains(&format!("server euid {server_uid}")),
2156 "reason did not name the rejected ACE: {reason}"
2157 );
2158 assert!(
2159 !data_root.join("config.json").exists(),
2160 "Haematite touched its ambient path before the UUID ACL refusal"
2161 );
2162 Ok(())
2163 }
2164
2165 #[cfg(all(feature = "haematite-backend", target_os = "macos"))]
2166 #[test]
2167 fn path_ambient_haematite_accepts_a_deny_only_acl_ancestor()
2168 -> Result<(), Box<dyn std::error::Error>> {
2169 use std::os::unix::fs::PermissionsExt as _;
2170
2171 let sandbox = crate::test_support::private_tempdir()?;
2172 std::fs::set_permissions(sandbox.path(), std::fs::Permissions::from_mode(0o700))?;
2173 let private_parent = sandbox.path().join("deny-only");
2174 let data_root = private_parent.join("data");
2175 std::fs::create_dir(&private_parent)?;
2176 std::fs::set_permissions(&private_parent, std::fs::Permissions::from_mode(0o700))?;
2177 let status = std::process::Command::new("chmod")
2178 .arg("+a")
2179 .arg("everyone deny delete")
2180 .arg(&private_parent)
2181 .status()?;
2182 assert!(status.success(), "failed to install Darwin deny-only ACL");
2183 let configured = data_root
2184 .to_str()
2185 .ok_or("temporary data path was not UTF-8")?;
2186
2187 let result = super::build_haematite_store(configured, 4, None);
2188 let cleanup = std::process::Command::new("chmod")
2189 .arg("-RN")
2190 .arg(&private_parent)
2191 .status()?;
2192 assert!(cleanup.success(), "failed to clean Darwin deny-only ACL");
2193
2194 let (store, responder) = result?;
2195 assert!(responder.is_none());
2196 assert!(data_root.join("config.json").is_file());
2197 drop(store);
2198 Ok(())
2199 }
2200
2201 #[cfg(all(feature = "haematite-backend", target_os = "macos"))]
2202 #[test]
2203 fn path_ambient_haematite_accepts_the_stock_home_acl_chain()
2204 -> Result<(), Box<dyn std::error::Error>> {
2205 use std::os::unix::fs::PermissionsExt as _;
2206 use users::os::unix::UserExt as _;
2207
2208 let effective_uid = rustix::process::geteuid().as_raw();
2209 let effective_user = users::get_user_by_uid(effective_uid)
2210 .ok_or_else(|| format!("server euid {effective_uid} has no account record"))?;
2211 let sandbox = tempfile::Builder::new()
2212 .prefix(".aion-acl-home-proof-")
2213 .tempdir_in(effective_user.home_dir())?;
2214 std::fs::set_permissions(sandbox.path(), std::fs::Permissions::from_mode(0o700))?;
2215 let data_root = sandbox.path().join("data");
2216 let configured = data_root
2217 .to_str()
2218 .ok_or("temporary data path was not UTF-8")?;
2219
2220 let (store, responder) = super::build_haematite_store(configured, 4, None)?;
2221 assert!(responder.is_none());
2222 assert!(data_root.join("config.json").is_file());
2223 drop(store);
2224 Ok(())
2225 }
2226
2227 #[cfg(all(
2228 feature = "haematite-backend",
2229 unix,
2230 not(any(target_os = "linux", target_os = "android"))
2231 ))]
2232 #[test]
2233 fn path_ambient_haematite_accepts_an_owner_controlled_chain()
2234 -> Result<(), Box<dyn std::error::Error>> {
2235 use std::os::unix::fs::PermissionsExt as _;
2236
2237 let sandbox = crate::test_support::private_tempdir()?;
2238 std::fs::set_permissions(sandbox.path(), std::fs::Permissions::from_mode(0o700))?;
2239 let private_parent = sandbox.path().join("private");
2240 let data_root = private_parent.join("data");
2241 std::fs::create_dir(&private_parent)?;
2242 std::fs::set_permissions(&private_parent, std::fs::Permissions::from_mode(0o700))?;
2243 let configured = data_root
2244 .to_str()
2245 .ok_or("temporary data path was not UTF-8")?;
2246
2247 let (store, responder) = super::build_haematite_store(configured, 4, None)?;
2248 assert!(responder.is_none());
2249 assert!(data_root.join("config.json").is_file());
2250 for shard in 0..4 {
2251 assert!(data_root.join(format!("shard-{shard}")).is_dir());
2252 }
2253 drop(store);
2254 Ok(())
2255 }
2256
2257 #[tokio::test]
2258 async fn connect_store_memory_backend_exposes_no_outbox_store()
2259 -> Result<(), Box<dyn std::error::Error>> {
2260 use crate::config::{StoreBackend, StoreConfig};
2261
2262 let connected = super::connect_store(StoreConfig {
2265 backend: StoreBackend::Memory,
2266 url: None,
2267 owned_shards: Vec::new(),
2268 data_dir: None,
2269 shard_count: 1,
2270 cluster: None,
2271 })
2272 .await?;
2273 assert!(
2274 connected.outbox_store.is_none(),
2275 "the in-memory backend exposes no outbox store"
2276 );
2277 Ok(())
2278 }
2279
2280 #[cfg(feature = "libsql-backend")]
2284 #[tokio::test]
2285 async fn connect_store_shares_outbox_store_only_for_libsql()
2286 -> Result<(), Box<dyn std::error::Error>> {
2287 use crate::config::{StoreBackend, StoreConfig};
2288
2289 let path = std::env::temp_dir().join(format!(
2294 "aion-connect-store-{}-{}.db",
2295 std::process::id(),
2296 std::time::SystemTime::now()
2297 .duration_since(std::time::UNIX_EPOCH)
2298 .map(|elapsed| elapsed.as_nanos())
2299 .unwrap_or_default()
2300 ));
2301 let connected = super::connect_store(StoreConfig {
2302 backend: StoreBackend::LibSql,
2303 url: Some(path.to_string_lossy().into_owned()),
2304 owned_shards: Vec::new(),
2305 data_dir: None,
2306 shard_count: 1,
2307 cluster: None,
2308 })
2309 .await?;
2310 assert!(
2311 connected.outbox_store.is_some(),
2312 "the libSQL backend shares its leaf store as the dispatcher's outbox store"
2313 );
2314 Ok(())
2315 }
2316
2317 #[tokio::test]
2318 async fn state_build_fails_without_event_broadcast_capacity()
2319 -> Result<(), Box<dyn std::error::Error>> {
2320 let mut runtime = runtime_config();
2321 runtime.websocket.event_broadcast_capacity = None;
2322
2323 let error = ServerState::build_with_store(InMemoryStore::default(), runtime)
2324 .await
2325 .err()
2326 .ok_or("state build must fail when event streaming is unsized")?;
2327
2328 assert!(error.is_config(), "expected a config error, got {error}");
2329 assert!(
2330 error
2331 .to_string()
2332 .contains("websocket.event_broadcast_capacity"),
2333 "error must name the missing key: {error}"
2334 );
2335 Ok(())
2336 }
2337
2338 #[tokio::test]
2339 async fn state_build_fails_without_query_timeout() -> Result<(), Box<dyn std::error::Error>> {
2340 let mut runtime = runtime_config();
2341 runtime.query_timeout = None;
2342
2343 let error = ServerState::build_with_store(InMemoryStore::default(), runtime)
2344 .await
2345 .err()
2346 .ok_or("state build must fail when the query reply deadline is unset")?;
2347
2348 assert!(error.is_config(), "expected a config error, got {error}");
2349 assert!(
2350 error.to_string().contains("runtime.query_timeout_ms"),
2351 "error must name the missing key: {error}"
2352 );
2353 assert!(
2354 error.to_string().contains("AION_RUNTIME_QUERY_TIMEOUT_MS"),
2355 "error must name the environment override: {error}"
2356 );
2357 Ok(())
2358 }
2359
2360 #[tokio::test]
2361 async fn state_build_fails_with_zero_query_timeout() -> Result<(), Box<dyn std::error::Error>> {
2362 let mut runtime = runtime_config();
2363 runtime.query_timeout = Some(Duration::ZERO);
2364
2365 let error = ServerState::build_with_store(InMemoryStore::default(), runtime)
2366 .await
2367 .err()
2368 .ok_or("state build must fail when the query reply deadline is zero")?;
2369
2370 assert!(error.is_config(), "expected a config error, got {error}");
2371 assert!(
2372 error.to_string().contains("runtime.query_timeout_ms"),
2373 "error must name the zero-valued key: {error}"
2374 );
2375 Ok(())
2376 }
2377}