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 queue_service_state: crate::worker::QueueServiceState,
94 queue_declarations: crate::worker::QueueDeclarationSource,
98 cluster_self_node: Option<String>,
103 #[cfg(feature = "haematite-backend")]
108 cluster_responder: Option<aion_store_haematite::ClusterResponder>,
109 #[cfg(feature = "haematite-backend")]
112 cluster_store: Option<Arc<aion_store_haematite::HaematiteStore>>,
113 #[cfg(feature = "haematite-backend")]
116 watched_peers: Vec<crate::cluster::WatchedPeer>,
117 #[cfg(feature = "haematite-backend")]
122 shard_directory: Option<Arc<crate::routing::StaticShardDirectory>>,
123 #[cfg(feature = "haematite-backend")]
127 request_forwarder: Option<Arc<dyn crate::routing::RequestForwarder>>,
128 #[cfg(feature = "auth")]
129 jwks_cache: Option<JwksCache>,
130}
131
132impl ServerState {
133 const FALLBACK_CLUSTER_BROADCAST_CAPACITY: std::num::NonZeroUsize =
142 match std::num::NonZeroUsize::new(64) {
143 Some(value) => value,
144 None => std::num::NonZeroUsize::MIN,
145 };
146
147 pub async fn build(config: ServerConfig) -> Result<Self, ServerError> {
154 let (store_config, runtime) = config.into_parts();
155 let connected = connect_store(store_config).await?;
156 Self::build_with_connected_store(connected, runtime).await
157 }
158
159 pub async fn build_with_store<S>(store: S, runtime: RuntimeConfig) -> Result<Self, ServerError>
165 where
166 S: EventStore + NamespaceStore,
167 {
168 let leaf = Arc::new(store);
172 let namespace_store: Arc<dyn NamespaceStore> = leaf.clone();
173 Self::build_with_connected_store(
174 ConnectedStore::local(leaf, None, namespace_store),
175 runtime,
176 )
177 .await
178 }
179
180 async fn build_with_connected_store(
181 connected: ConnectedStore,
182 runtime: RuntimeConfig,
183 ) -> Result<Self, ServerError> {
184 let outbox_store = connected.outbox_store;
185 let bootstrap_coordinator = connected.bootstrap_coordinator;
186 #[cfg(feature = "haematite-backend")]
187 let cluster_responder = connected.cluster_responder;
188 #[cfg(feature = "haematite-backend")]
189 let cluster_store = connected.cluster_store;
190 #[cfg(feature = "haematite-backend")]
191 let watched_peers = connected.watched_peers;
192 #[cfg(feature = "haematite-backend")]
195 let cluster_self_node = connected.self_node_id.clone();
196 #[cfg(not(feature = "haematite-backend"))]
197 let cluster_self_node: Option<String> = None;
198 #[cfg(feature = "haematite-backend")]
202 let RoutingState {
203 shard_directory,
204 request_forwarder,
205 mint_routing,
206 } = build_routing_state(
207 cluster_store.as_ref(),
208 connected.directory_peers,
209 connected.self_node_id,
210 );
211 let (event_broadcast_capacity, query_timeout) = required_engine_seams(&runtime)?;
212 let (cluster_publisher, transcript_publisher) =
213 build_real_time_publishers(&runtime, connected.observability_store)?;
214 let metrics = Metrics::new().map_err(|error| metrics_config_error(&error))?;
215 let outbox_wake = Arc::new(tokio::sync::Notify::new());
220 let instrumented_store = Arc::new(
221 InstrumentedEventStore::new(
222 connected.event_store,
223 metrics.clone(),
224 runtime.default_namespace.clone(),
225 )
226 .with_outbox_wake(Arc::clone(&outbox_wake)),
227 );
228 let exported_metrics = runtime.metrics.enabled.then_some(metrics.clone());
229 #[cfg(not(feature = "haematite-backend"))]
230 let mint_routing: Option<crate::namespace::NamespaceRouting> = None;
231 let seams = build_worker_seams(
232 &runtime,
233 &cluster_publisher,
234 &connected.namespace_store,
235 mint_routing,
236 );
237 let (activity_dispatcher, activity_mock_registry, attempt_owners) =
238 build_decorated_dispatcher(&runtime, &seams);
239
240 let engine = build_engine(EngineAssembly {
241 instrumented_store: &instrumented_store,
242 event_broadcast_capacity,
243 query_timeout,
244 activity_dispatcher,
245 active_registry: Arc::new(aion::Registry::default()),
246 bootstrap_coordinator,
247 runtime: &runtime,
248 })
249 .await?;
250 let engine = Arc::new(engine);
251 install_engine_backed_seams(&seams, &engine, runtime.outbox.enabled);
252 let resolver = NamespaceResolver::from_config(runtime.namespace.clone(), engine);
253 #[cfg(feature = "auth")]
254 let jwks_cache = build_jwks_cache(&runtime).await?;
255 Ok(Self {
256 inner: Arc::new(ServerStateInner {
257 namespace_guard: NamespaceGuard::new(resolver),
258 runtime,
259 worker_registry: seams.worker_registry,
260 pending_activities: seams.pending_activities,
261 heartbeat_tracker: seams.heartbeat_tracker,
262 drain_state: seams.drain_state,
263 metrics: exported_metrics,
264 health: Some(HealthState::new(instrumented_store, true)),
265 activity_mock_registry,
266 outbox_store,
267 namespace_store: connected.namespace_store,
268 outbox_wake,
269 cluster_publisher,
270 transcript_publisher,
271 attempt_owners,
272 queue_service_state: seams.queue_service_state,
273 queue_declarations: seams.queue_declarations,
274 cluster_self_node,
275 #[cfg(feature = "haematite-backend")]
276 cluster_responder,
277 #[cfg(feature = "haematite-backend")]
278 cluster_store,
279 #[cfg(feature = "haematite-backend")]
280 watched_peers,
281 #[cfg(feature = "haematite-backend")]
282 shard_directory,
283 #[cfg(feature = "haematite-backend")]
284 request_forwarder,
285 #[cfg(feature = "auth")]
286 jwks_cache,
287 }),
288 })
289 }
290
291 #[must_use]
293 pub fn from_parts(namespace_resolver: NamespaceResolver, runtime: RuntimeConfig) -> Self {
294 Self::from_parts_with_namespace_store(
299 namespace_resolver,
300 runtime,
301 Arc::new(aion_store::InMemoryStore::default()),
302 )
303 }
304
305 #[must_use]
313 pub fn from_parts_with_namespace_store(
314 namespace_resolver: NamespaceResolver,
315 runtime: RuntimeConfig,
316 namespace_store: Arc<dyn NamespaceStore>,
317 ) -> Self {
318 let heartbeat_tracker = HeartbeatTracker::new(runtime.worker.heartbeat_window);
319 let pending_activities =
326 PendingActivities::default().with_heartbeat_window(runtime.worker.heartbeat_window);
327 let bounds = transcript_bounds(&runtime);
331 Self {
332 inner: Arc::new(ServerStateInner {
333 namespace_guard: NamespaceGuard::new(namespace_resolver),
334 runtime,
335 worker_registry: ConnectedWorkerRegistry::default(),
336 pending_activities,
337 heartbeat_tracker,
338 drain_state: DrainState::default(),
339 metrics: None,
340 health: None,
341 activity_mock_registry: None,
342 outbox_store: None,
343 namespace_store,
344 outbox_wake: Arc::new(tokio::sync::Notify::new()),
345 cluster_publisher: crate::cluster_publisher::ClusterEventPublisher::new(
346 Self::FALLBACK_CLUSTER_BROADCAST_CAPACITY,
347 ),
348 transcript_publisher: build_transcript_publisher(
352 None,
353 Self::FALLBACK_CLUSTER_BROADCAST_CAPACITY,
354 bounds,
355 ),
356 attempt_owners: crate::worker::AttemptOwnerIndex::new(),
357 queue_service_state: crate::worker::QueueServiceState::default(),
358 queue_declarations: crate::worker::QueueDeclarationSource::default(),
359 cluster_self_node: None,
360 #[cfg(feature = "haematite-backend")]
361 cluster_responder: None,
362 #[cfg(feature = "haematite-backend")]
363 cluster_store: None,
364 #[cfg(feature = "haematite-backend")]
365 watched_peers: Vec::new(),
366 #[cfg(feature = "haematite-backend")]
367 shard_directory: None,
368 #[cfg(feature = "haematite-backend")]
369 request_forwarder: None,
370 #[cfg(feature = "auth")]
371 jwks_cache: None,
372 }),
373 }
374 }
375
376 #[cfg(feature = "auth")]
385 #[must_use]
386 pub fn from_parts_with_namespace_store_and_jwks(
387 namespace_resolver: NamespaceResolver,
388 runtime: RuntimeConfig,
389 namespace_store: Arc<dyn NamespaceStore>,
390 jwks_cache: JwksCache,
391 ) -> Self {
392 let heartbeat_tracker = HeartbeatTracker::new(runtime.worker.heartbeat_window);
393 let pending_activities =
400 PendingActivities::default().with_heartbeat_window(runtime.worker.heartbeat_window);
401 let bounds = transcript_bounds(&runtime);
405 Self {
406 inner: Arc::new(ServerStateInner {
407 namespace_guard: NamespaceGuard::new(namespace_resolver),
408 runtime,
409 worker_registry: ConnectedWorkerRegistry::default(),
410 pending_activities,
411 heartbeat_tracker,
412 drain_state: DrainState::default(),
413 metrics: None,
414 health: None,
415 activity_mock_registry: None,
416 outbox_store: None,
417 namespace_store,
418 outbox_wake: Arc::new(tokio::sync::Notify::new()),
419 cluster_publisher: crate::cluster_publisher::ClusterEventPublisher::new(
420 Self::FALLBACK_CLUSTER_BROADCAST_CAPACITY,
421 ),
422 transcript_publisher: build_transcript_publisher(
426 None,
427 Self::FALLBACK_CLUSTER_BROADCAST_CAPACITY,
428 bounds,
429 ),
430 attempt_owners: crate::worker::AttemptOwnerIndex::new(),
431 queue_service_state: crate::worker::QueueServiceState::default(),
432 queue_declarations: crate::worker::QueueDeclarationSource::default(),
433 cluster_self_node: None,
434 #[cfg(feature = "haematite-backend")]
435 cluster_responder: None,
436 #[cfg(feature = "haematite-backend")]
437 cluster_store: None,
438 #[cfg(feature = "haematite-backend")]
439 watched_peers: Vec::new(),
440 #[cfg(feature = "haematite-backend")]
441 shard_directory: None,
442 #[cfg(feature = "haematite-backend")]
443 request_forwarder: None,
444 jwks_cache: Some(jwks_cache),
445 }),
446 }
447 }
448
449 #[cfg(feature = "auth")]
455 #[must_use]
456 pub fn from_parts_with_jwks(
457 namespace_resolver: NamespaceResolver,
458 runtime: RuntimeConfig,
459 jwks_cache: JwksCache,
460 ) -> Self {
461 let heartbeat_tracker = HeartbeatTracker::new(runtime.worker.heartbeat_window);
462 let pending_activities =
469 PendingActivities::default().with_heartbeat_window(runtime.worker.heartbeat_window);
470 let bounds = transcript_bounds(&runtime);
474 Self {
475 inner: Arc::new(ServerStateInner {
476 namespace_guard: NamespaceGuard::new(namespace_resolver),
477 runtime,
478 worker_registry: ConnectedWorkerRegistry::default(),
479 pending_activities,
480 heartbeat_tracker,
481 drain_state: DrainState::default(),
482 metrics: None,
483 health: None,
484 activity_mock_registry: None,
485 outbox_store: None,
486 namespace_store: Arc::new(aion_store::InMemoryStore::default()),
491 outbox_wake: Arc::new(tokio::sync::Notify::new()),
492 cluster_publisher: crate::cluster_publisher::ClusterEventPublisher::new(
493 Self::FALLBACK_CLUSTER_BROADCAST_CAPACITY,
494 ),
495 transcript_publisher: build_transcript_publisher(
499 None,
500 Self::FALLBACK_CLUSTER_BROADCAST_CAPACITY,
501 bounds,
502 ),
503 attempt_owners: crate::worker::AttemptOwnerIndex::new(),
504 queue_service_state: crate::worker::QueueServiceState::default(),
505 queue_declarations: crate::worker::QueueDeclarationSource::default(),
506 cluster_self_node: None,
507 #[cfg(feature = "haematite-backend")]
508 cluster_responder: None,
509 #[cfg(feature = "haematite-backend")]
510 cluster_store: None,
511 #[cfg(feature = "haematite-backend")]
512 watched_peers: Vec::new(),
513 #[cfg(feature = "haematite-backend")]
514 shard_directory: None,
515 #[cfg(feature = "haematite-backend")]
516 request_forwarder: None,
517 jwks_cache: Some(jwks_cache),
518 }),
519 }
520 }
521
522 #[must_use]
524 pub fn from_parts_with_registry(
525 namespace_resolver: NamespaceResolver,
526 runtime: RuntimeConfig,
527 worker_registry: ConnectedWorkerRegistry,
528 ) -> Self {
529 let heartbeat_tracker = HeartbeatTracker::new(runtime.worker.heartbeat_window);
530 let pending_activities =
537 PendingActivities::default().with_heartbeat_window(runtime.worker.heartbeat_window);
538 let bounds = transcript_bounds(&runtime);
542 Self {
543 inner: Arc::new(ServerStateInner {
544 namespace_guard: NamespaceGuard::new(namespace_resolver),
545 runtime,
546 worker_registry,
547 pending_activities,
548 heartbeat_tracker,
549 drain_state: DrainState::default(),
550 metrics: None,
551 health: None,
552 activity_mock_registry: None,
553 outbox_store: None,
554 namespace_store: Arc::new(aion_store::InMemoryStore::default()),
559 outbox_wake: Arc::new(tokio::sync::Notify::new()),
560 cluster_publisher: crate::cluster_publisher::ClusterEventPublisher::new(
561 Self::FALLBACK_CLUSTER_BROADCAST_CAPACITY,
562 ),
563 transcript_publisher: build_transcript_publisher(
567 None,
568 Self::FALLBACK_CLUSTER_BROADCAST_CAPACITY,
569 bounds,
570 ),
571 attempt_owners: crate::worker::AttemptOwnerIndex::new(),
572 queue_service_state: crate::worker::QueueServiceState::default(),
573 queue_declarations: crate::worker::QueueDeclarationSource::default(),
574 cluster_self_node: None,
575 #[cfg(feature = "haematite-backend")]
576 cluster_responder: None,
577 #[cfg(feature = "haematite-backend")]
578 cluster_store: None,
579 #[cfg(feature = "haematite-backend")]
580 watched_peers: Vec::new(),
581 #[cfg(feature = "haematite-backend")]
582 shard_directory: None,
583 #[cfg(feature = "haematite-backend")]
584 request_forwarder: None,
585 #[cfg(feature = "auth")]
586 jwks_cache: None,
587 }),
588 }
589 }
590
591 #[must_use]
593 pub fn namespace_guard(&self) -> &NamespaceGuard {
594 &self.inner.namespace_guard
595 }
596
597 #[must_use]
599 pub fn deploy_guard(&self) -> crate::deploy::DeployGuard {
600 crate::deploy::DeployGuard::new(self.inner.namespace_guard.resolver().clone())
601 }
602
603 #[must_use]
605 pub fn runtime_config(&self) -> &RuntimeConfig {
606 &self.inner.runtime
607 }
608
609 #[must_use]
611 pub fn worker_registry(&self) -> &ConnectedWorkerRegistry {
612 &self.inner.worker_registry
613 }
614
615 #[must_use]
619 pub fn cluster_publisher(&self) -> &crate::cluster_publisher::ClusterEventPublisher {
620 &self.inner.cluster_publisher
621 }
622
623 #[must_use]
628 pub fn transcript_publisher(&self) -> &crate::activity_publisher::ActivityEventPublisher {
629 &self.inner.transcript_publisher
630 }
631
632 #[must_use]
636 pub fn attempt_owners(&self) -> &crate::worker::AttemptOwnerIndex {
637 &self.inner.attempt_owners
638 }
639
640 #[must_use]
643 pub fn queue_service_state(&self) -> &crate::worker::QueueServiceState {
644 &self.inner.queue_service_state
645 }
646
647 #[must_use]
651 pub fn queue_declarations(&self) -> &crate::worker::QueueDeclarationSource {
652 &self.inner.queue_declarations
653 }
654
655 pub fn unserved_queues(&self) -> Result<Vec<crate::worker::UnservedQueue>, ServerError> {
666 self.inner.queue_service_state.unserved()
667 }
668
669 pub fn unrecoverable_runs(
687 &self,
688 ) -> Result<Vec<(aion_core::WorkflowId, aion::registry::UnrecoverableRun)>, ServerError> {
689 self.engine()?
690 .registry()
691 .unrecoverable()
692 .list()
693 .map_err(ServerError::from)
694 }
695
696 #[must_use]
708 pub fn intervention_router(&self) -> crate::worker::InterventionRouter {
709 let transport: std::sync::Arc<dyn crate::worker::InterventionTransport> = {
710 #[cfg(feature = "liminal-transport")]
711 {
712 std::sync::Arc::new(crate::worker::LiminalInterventionTransport)
713 }
714 #[cfg(not(feature = "liminal-transport"))]
715 {
716 std::sync::Arc::new(NullInterventionTransport)
717 }
718 };
719 crate::worker::InterventionRouter::new(
720 self.inner.worker_registry.clone(),
721 self.inner.attempt_owners.clone(),
722 transport,
723 )
724 .with_transcript_publisher(self.inner.transcript_publisher.clone())
727 }
728
729 #[must_use]
733 pub fn cluster_self_node(&self) -> Option<&str> {
734 self.inner.cluster_self_node.as_deref()
735 }
736
737 pub fn engine(&self) -> Result<Arc<aion::Engine>, ServerError> {
750 self.inner
751 .namespace_guard
752 .resolver()
753 .engine()
754 .map(Arc::clone)
755 }
756
757 #[must_use]
759 pub fn pending_activities(&self) -> &PendingActivities {
760 &self.inner.pending_activities
761 }
762
763 #[must_use]
765 pub fn heartbeat_tracker(&self) -> &HeartbeatTracker {
766 &self.inner.heartbeat_tracker
767 }
768
769 #[must_use]
771 pub fn drain_state(&self) -> &DrainState {
772 &self.inner.drain_state
773 }
774
775 #[must_use]
777 pub fn metrics(&self) -> Option<&Metrics> {
778 self.inner.metrics.as_ref()
779 }
780
781 #[must_use]
783 pub fn health(&self) -> Option<&HealthState> {
784 self.inner.health.as_ref()
785 }
786
787 #[must_use]
792 pub fn activity_mock_registry(&self) -> Option<&ActivityMockRegistry> {
793 self.inner.activity_mock_registry.as_ref()
794 }
795
796 #[must_use]
802 pub fn outbox_store(&self) -> Option<Arc<dyn OutboxStore>> {
803 self.inner.outbox_store.clone()
804 }
805
806 #[must_use]
815 pub fn namespace_store(&self) -> &Arc<dyn NamespaceStore> {
816 &self.inner.namespace_store
817 }
818
819 #[must_use]
828 pub fn namespace_minter(&self) -> NamespaceMinter {
829 let minter = NamespaceMinter::new(
830 Arc::clone(&self.inner.namespace_store),
831 self.inner.runtime.auto_create,
832 )
833 .with_cluster_publisher(self.inner.cluster_publisher.clone());
838 match self.namespace_routing() {
843 Some(routing) => minter.with_routing(routing),
844 None => minter,
845 }
846 }
847
848 #[cfg(feature = "haematite-backend")]
858 #[must_use]
859 pub fn namespace_routing(&self) -> Option<crate::namespace::NamespaceRouting> {
860 build_namespace_routing(
861 self.cluster_store(),
862 self.shard_directory(),
863 self.request_forwarder(),
864 )
865 }
866
867 #[cfg(not(feature = "haematite-backend"))]
869 #[must_use]
870 pub const fn namespace_routing(&self) -> Option<crate::namespace::NamespaceRouting> {
871 None
872 }
873
874 #[must_use]
880 pub fn outbox_wake(&self) -> Arc<tokio::sync::Notify> {
881 Arc::clone(&self.inner.outbox_wake)
882 }
883
884 #[cfg(feature = "haematite-backend")]
890 #[must_use]
891 pub fn is_clustered(&self) -> bool {
892 self.inner.cluster_responder.is_some()
893 }
894
895 #[cfg(feature = "haematite-backend")]
900 #[must_use]
901 pub fn cluster_store(&self) -> Option<&Arc<aion_store_haematite::HaematiteStore>> {
902 self.inner.cluster_store.as_ref()
903 }
904
905 #[cfg(feature = "haematite-backend")]
909 #[must_use]
910 pub fn shard_directory(&self) -> Option<&Arc<crate::routing::StaticShardDirectory>> {
911 self.inner.shard_directory.as_ref()
912 }
913
914 #[cfg(feature = "haematite-backend")]
917 #[must_use]
918 pub fn request_forwarder(&self) -> Option<&Arc<dyn crate::routing::RequestForwarder>> {
919 self.inner.request_forwarder.as_ref()
920 }
921
922 #[must_use]
939 pub fn spawn_heartbeat_sweeper(
940 &self,
941 shutdown: tokio::sync::watch::Receiver<bool>,
942 ) -> tokio::task::JoinHandle<()> {
943 let sweeper = crate::worker::HeartbeatSweeper::new(
944 self.inner.heartbeat_tracker.clone(),
945 self.inner.worker_registry.clone(),
946 self.inner.pending_activities.clone(),
947 self.inner.drain_state.clone(),
948 self.inner.runtime.worker.heartbeat_window,
949 )
950 .with_queue_state(self.inner.queue_service_state.clone());
954 tokio::spawn(sweeper.run(shutdown))
955 }
956
957 #[cfg(feature = "liminal-transport")]
968 #[must_use]
969 pub fn spawn_liminal_liveness_probe(
970 &self,
971 notifier: std::sync::Arc<crate::worker::LiminalConnectionNotifier>,
972 shutdown: tokio::sync::watch::Receiver<bool>,
973 ) -> tokio::task::JoinHandle<()> {
974 let probe = crate::worker::LivenessProbe::new(
975 notifier,
976 self.inner.heartbeat_tracker.clone(),
977 self.inner.worker_registry.clone(),
978 self.inner.runtime.worker.heartbeat_window,
979 );
980 tokio::spawn(probe.run(shutdown))
981 }
982
983 #[cfg(feature = "haematite-backend")]
999 pub fn spawn_cluster_supervisor(
1000 &self,
1001 config: crate::cluster::SupervisorConfig,
1002 shutdown: tokio::sync::watch::Receiver<bool>,
1003 ) -> Result<bool, ServerError> {
1004 let Some(cluster_store) = self.inner.cluster_store.clone() else {
1005 return Ok(false);
1006 };
1007 if self.inner.watched_peers.is_empty() {
1008 return Ok(false);
1009 }
1010 let engine = Arc::clone(self.inner.namespace_guard.resolver().engine()?);
1011 let publisher = Arc::new(self.inner.cluster_publisher.clone());
1015 let self_node = self.inner.cluster_self_node.clone().unwrap_or_default();
1016 let adopter = Arc::new(crate::cluster::OutboxSettlingAdopter::new(
1022 engine,
1023 self.inner.outbox_store.clone(),
1024 ));
1025 let supervisor = crate::cluster::ClusterSupervisor::new(
1026 cluster_store,
1027 adopter,
1028 self.inner.watched_peers.clone(),
1029 config,
1030 )
1031 .with_publisher(publisher, self_node);
1032 if !supervisor.watches_any() {
1033 return Ok(false);
1034 }
1035 tokio::spawn(supervisor.run(shutdown));
1036 Ok(true)
1037 }
1038
1039 #[cfg(feature = "auth")]
1041 #[must_use]
1042 pub fn jwks_cache(&self) -> Option<&JwksCache> {
1043 self.inner.jwks_cache.as_ref()
1044 }
1045
1046 pub fn shutdown(&self) -> Result<(), ServerError> {
1053 self.inner.namespace_guard.resolver().shutdown_engine()
1054 }
1055}
1056
1057#[cfg(feature = "auth")]
1058async fn build_jwks_cache(runtime: &RuntimeConfig) -> Result<Option<JwksCache>, ServerError> {
1059 if !runtime.auth.enabled {
1060 return Ok(None);
1061 }
1062 let Some(url) = runtime.auth.jwks_url.clone() else {
1063 return Err(ServerError::Config {
1064 message: "auth.jwks_url must not be empty when auth.enabled is true".to_owned(),
1065 });
1066 };
1067 let interval = std::time::Duration::from_secs(runtime.auth.jwks_refresh_seconds);
1068 let cache = JwksCache::new(url, interval)
1069 .await
1070 .map_err(|error| ServerError::Config {
1071 message: format!("auth jwks initial fetch failed: {error}"),
1072 })?;
1073 Ok(Some(cache))
1074}
1075
1076fn metrics_config_error(error: &MetricsError) -> ServerError {
1077 ServerError::Config {
1078 message: error.to_string(),
1079 }
1080}
1081
1082struct EngineAssembly<'a> {
1084 instrumented_store: &'a Arc<InstrumentedEventStore>,
1086 event_broadcast_capacity: std::num::NonZeroUsize,
1088 query_timeout: std::time::Duration,
1090 activity_dispatcher: Arc<dyn ActivityDispatcher>,
1092 active_registry: Arc<aion::Registry>,
1094 bootstrap_coordinator: bool,
1096 runtime: &'a RuntimeConfig,
1098}
1099
1100async fn build_engine(assembly: EngineAssembly<'_>) -> Result<aion::Engine, ServerError> {
1107 let mut search_attribute_schema = aion_core::SearchAttributeSchema::new();
1108 search_attribute_schema
1109 .register(
1110 crate::namespace::NAMESPACE_ATTRIBUTE,
1111 aion_core::SearchAttributeType::String,
1112 )
1113 .map_err(|error| ServerError::Config {
1114 message: format!("failed to register namespace search attribute: {error}"),
1115 })?;
1116 search_attribute_schema
1117 .register(
1118 crate::namespace::TASK_QUEUE_ATTRIBUTE,
1119 aion_core::SearchAttributeType::String,
1120 )
1121 .map_err(|error| ServerError::Config {
1122 message: format!("failed to register task_queue search attribute: {error}"),
1123 })?;
1124 let runtime = assembly.runtime;
1125 let builder = EngineBuilder::new()
1126 .store_arc(assembly.instrumented_store.clone())
1127 .event_streaming(assembly.event_broadcast_capacity)
1128 .in_memory_visibility()
1129 .search_attribute_schema(search_attribute_schema)
1130 .scheduler_threads(runtime.scheduler_threads)
1131 .outbox_enabled(runtime.outbox.enabled)
1132 .activity_dispatcher(assembly.activity_dispatcher)
1133 .active_registry(assembly.active_registry)
1134 .production_recovery_seam()
1135 .signal_router_factory(|runtime: Arc<RuntimeHandle>, handoff| {
1136 Arc::new(ConcreteSignalRouter::new(runtime, handoff)) as Arc<dyn SignalRouter>
1137 })
1138 .query_timeout(assembly.query_timeout)
1139 .bootstrap_schedule_coordinator(assembly.bootstrap_coordinator)
1145 .load_workflow_sources(runtime.workflow_packages.iter().map(PathBuf::as_path));
1146 let builder = if runtime.owned_shards.is_empty() {
1152 builder
1153 } else {
1154 builder.owned_shards(runtime.owned_shards.iter().copied())
1155 };
1156 builder.build().await.map_err(ServerError::from)
1157}
1158
1159fn required_engine_seams(
1164 runtime: &RuntimeConfig,
1165) -> Result<(std::num::NonZeroUsize, std::time::Duration), ServerError> {
1166 let event_broadcast_capacity = runtime
1167 .websocket
1168 .event_broadcast_capacity
1169 .and_then(std::num::NonZeroUsize::new)
1170 .ok_or_else(|| ServerError::Config {
1171 message: crate::config::EVENT_BROADCAST_CAPACITY_REQUIRED.to_owned(),
1172 })?;
1173 let query_timeout = runtime
1174 .query_timeout
1175 .filter(|timeout| !timeout.is_zero())
1176 .ok_or_else(|| ServerError::Config {
1177 message: crate::config::QUERY_TIMEOUT_REQUIRED.to_owned(),
1178 })?;
1179 Ok((event_broadcast_capacity, query_timeout))
1180}
1181
1182fn install_outbox_delivery(
1189 pending_activities: &PendingActivities,
1190 engine: &Arc<aion::Engine>,
1191 outbox_enabled: bool,
1192) {
1193 if outbox_enabled {
1194 let callback = Arc::new(crate::worker::ServerOutboxDeliveryCallback::new(
1195 Arc::clone(engine),
1196 ));
1197 pending_activities.set_outbox_delivery(callback);
1198 }
1199}
1200
1201struct WorkerSeams {
1203 worker_registry: ConnectedWorkerRegistry,
1204 pending_activities: PendingActivities,
1205 heartbeat_tracker: HeartbeatTracker,
1206 drain_state: DrainState,
1207 queue_declarations: crate::worker::QueueDeclarationSource,
1208 queue_service_state: crate::worker::QueueServiceState,
1209 declared_bodies: crate::worker::DeclaredBodySource,
1210}
1211
1212fn build_worker_seams(
1217 runtime: &RuntimeConfig,
1218 cluster_publisher: &crate::cluster_publisher::ClusterEventPublisher,
1219 namespace_store: &Arc<dyn NamespaceStore>,
1220 mint_routing: Option<crate::namespace::NamespaceRouting>,
1221) -> WorkerSeams {
1222 let worker_registry = ConnectedWorkerRegistry::default()
1223 .with_cluster_publisher(cluster_publisher.clone())
1224 .with_namespace_minting(namespace_store.clone(), runtime.auto_create);
1225 let worker_registry = match mint_routing {
1231 Some(routing) => worker_registry.with_namespace_routing(routing),
1232 None => worker_registry,
1233 };
1234 WorkerSeams {
1235 worker_registry,
1236 pending_activities: PendingActivities::default()
1241 .with_heartbeat_window(runtime.worker.heartbeat_window),
1242 heartbeat_tracker: HeartbeatTracker::new(runtime.worker.heartbeat_window),
1243 drain_state: DrainState::default(),
1244 queue_declarations: crate::worker::QueueDeclarationSource::default(),
1245 queue_service_state: crate::worker::QueueServiceState::default(),
1246 declared_bodies: crate::worker::DeclaredBodySource::default(),
1247 }
1248}
1249
1250fn install_queue_declarations(
1257 queue_declarations: &crate::worker::QueueDeclarationSource,
1258 engine: &Arc<aion::Engine>,
1259) {
1260 queue_declarations.install(Arc::new(crate::worker::EngineQueueDeclarations::new(
1261 Arc::clone(engine),
1262 )));
1263}
1264
1265fn install_declared_bodies(
1273 declared_bodies: &crate::worker::DeclaredBodySource,
1274 engine: &Arc<aion::Engine>,
1275) {
1276 declared_bodies.install(Arc::new(crate::worker::EngineDeclaredBodies::new(
1277 Arc::clone(engine),
1278 )));
1279}
1280
1281fn install_engine_backed_seams(seams: &WorkerSeams, engine: &Arc<aion::Engine>, outbox: bool) {
1287 install_outbox_delivery(&seams.pending_activities, engine, outbox);
1288 install_queue_declarations(&seams.queue_declarations, engine);
1289 install_declared_bodies(&seams.declared_bodies, engine);
1290}
1291
1292fn build_bridge_dispatcher(
1305 runtime: &RuntimeConfig,
1306 seams: &WorkerSeams,
1307) -> (WorkerActivityDispatcher, crate::worker::AttemptOwnerIndex) {
1308 let attempt_owners = crate::worker::AttemptOwnerIndex::new();
1309 let dispatcher = WorkerActivityDispatcher::new(
1310 seams.worker_registry.clone(),
1311 runtime.default_namespace.clone(),
1312 seams.heartbeat_tracker.clone(),
1313 )
1314 .with_pending(seams.pending_activities.clone())
1315 .with_drain_state(seams.drain_state.clone())
1316 .with_tokio_handle(tokio::runtime::Handle::current())
1317 .with_attempt_owners(attempt_owners.clone())
1318 .with_queue_service(runtime.worker.queue_service.clone())
1319 .with_queue_declarations(seams.queue_declarations.clone())
1320 .with_queue_state(seams.queue_service_state.clone());
1321 (dispatcher, attempt_owners)
1322}
1323
1324fn build_decorated_dispatcher(
1327 runtime: &RuntimeConfig,
1328 seams: &WorkerSeams,
1329) -> (
1330 Arc<dyn ActivityDispatcher>,
1331 Option<ActivityMockRegistry>,
1332 crate::worker::AttemptOwnerIndex,
1333) {
1334 let (dispatcher, attempt_owners) = build_bridge_dispatcher(runtime, seams);
1335 let (activity_dispatcher, activity_mock_registry) = decorate_activity_dispatcher(
1336 dispatcher,
1337 seams.declared_bodies.clone(),
1338 runtime.dev.enabled,
1339 );
1340 (activity_dispatcher, activity_mock_registry, attempt_owners)
1341}
1342
1343fn decorate_activity_dispatcher(
1344 dispatcher: WorkerActivityDispatcher,
1345 declared_bodies: crate::worker::DeclaredBodySource,
1346 dev_enabled: bool,
1347) -> (Arc<dyn ActivityDispatcher>, Option<ActivityMockRegistry>) {
1348 let declared = crate::worker::DeclaredCommandDispatcher::new(
1354 Arc::new(dispatcher),
1355 declared_bodies,
1356 tokio::runtime::Handle::current(),
1357 );
1358 if dev_enabled {
1359 let registry = ActivityMockRegistry::new();
1360 let decorated = DevMockingDispatcher::new(Arc::new(declared), registry.clone());
1361 (Arc::new(decorated), Some(registry))
1362 } else {
1363 (Arc::new(declared), None)
1364 }
1365}
1366
1367fn required_cluster_broadcast_capacity(
1372 runtime: &RuntimeConfig,
1373) -> Result<std::num::NonZeroUsize, ServerError> {
1374 runtime
1375 .websocket
1376 .cluster_broadcast_capacity
1377 .and_then(std::num::NonZeroUsize::new)
1378 .ok_or_else(|| ServerError::Config {
1379 message: crate::config::CLUSTER_BROADCAST_CAPACITY_REQUIRED.to_owned(),
1380 })
1381}
1382
1383fn build_real_time_publishers(
1396 runtime: &RuntimeConfig,
1397 observability_store: Option<Arc<dyn aion_store::ObservabilityStore>>,
1398) -> Result<
1399 (
1400 crate::cluster_publisher::ClusterEventPublisher,
1401 crate::activity_publisher::ActivityEventPublisher,
1402 ),
1403 ServerError,
1404> {
1405 let capacity = required_cluster_broadcast_capacity(runtime)?;
1406 Ok((
1407 crate::cluster_publisher::ClusterEventPublisher::new(capacity),
1408 build_transcript_publisher(observability_store, capacity, transcript_bounds(runtime)),
1409 ))
1410}
1411
1412fn transcript_bounds(runtime: &RuntimeConfig) -> crate::activity_bounds::TranscriptBounds {
1414 crate::activity_bounds::TranscriptBounds {
1415 max_event_bytes: runtime.observability.max_event_bytes,
1416 max_stream_events: runtime.observability.max_stream_events,
1417 }
1418}
1419
1420fn build_transcript_publisher(
1431 observability_store: Option<Arc<dyn aion_store::ObservabilityStore>>,
1432 capacity: std::num::NonZeroUsize,
1433 bounds: crate::activity_bounds::TranscriptBounds,
1434) -> crate::activity_publisher::ActivityEventPublisher {
1435 let store = observability_store
1436 .unwrap_or_else(|| Arc::new(aion_store::InMemoryObservabilityStore::default()));
1437 crate::activity_publisher::ActivityEventPublisher::new(store, capacity).with_bounds(bounds)
1438}
1439
1440#[cfg(feature = "haematite-backend")]
1442struct RoutingState {
1443 shard_directory: Option<Arc<crate::routing::StaticShardDirectory>>,
1444 request_forwarder: Option<Arc<dyn crate::routing::RequestForwarder>>,
1445 mint_routing: Option<crate::namespace::NamespaceRouting>,
1450}
1451
1452#[cfg(feature = "haematite-backend")]
1456fn build_routing_state(
1457 cluster_store: Option<&Arc<aion_store_haematite::HaematiteStore>>,
1458 directory_peers: Vec<crate::routing::DirectoryPeer>,
1459 self_node_id: Option<String>,
1460) -> RoutingState {
1461 let Some(store) = cluster_store else {
1462 return RoutingState {
1463 shard_directory: None,
1464 request_forwarder: None,
1465 mint_routing: None,
1466 };
1467 };
1468 let shard_directory = Arc::new(crate::routing::StaticShardDirectory::new(
1469 Arc::clone(store),
1470 directory_peers,
1471 self_node_id,
1472 ));
1473 let request_forwarder: Arc<dyn crate::routing::RequestForwarder> =
1474 Arc::new(crate::routing::GrpcRequestForwarder::new());
1475 let mint_routing = build_namespace_routing(
1476 Some(store),
1477 Some(&shard_directory),
1478 Some(&request_forwarder),
1479 );
1480 RoutingState {
1481 shard_directory: Some(shard_directory),
1482 request_forwarder: Some(request_forwarder),
1483 mint_routing,
1484 }
1485}
1486
1487#[cfg(feature = "haematite-backend")]
1497fn build_namespace_routing(
1498 cluster_store: Option<&Arc<aion_store_haematite::HaematiteStore>>,
1499 shard_directory: Option<&Arc<crate::routing::StaticShardDirectory>>,
1500 request_forwarder: Option<&Arc<dyn crate::routing::RequestForwarder>>,
1501) -> Option<crate::namespace::NamespaceRouting> {
1502 use crate::namespace::{
1503 GrpcMintForwarder, MintForwarder, MintShardOwners, NamespaceRouting, NamespaceShardResolver,
1504 };
1505 let store = Arc::clone(cluster_store?);
1506 let directory = Arc::clone(shard_directory?);
1507 let shards: Arc<dyn NamespaceShardResolver> = store;
1508 let owners: Arc<dyn MintShardOwners> = directory;
1509 let forwarder: Arc<dyn MintForwarder> =
1510 Arc::new(GrpcMintForwarder::new(Arc::clone(request_forwarder?)));
1511 Some(NamespaceRouting::new(shards, owners, forwarder))
1512}
1513
1514struct ConnectedStore {
1524 event_store: Arc<dyn EventStore>,
1525 outbox_store: Option<Arc<dyn OutboxStore>>,
1526 namespace_store: Arc<dyn NamespaceStore>,
1533 observability_store: Option<Arc<dyn aion_store::ObservabilityStore>>,
1541 bootstrap_coordinator: bool,
1542 #[cfg(feature = "haematite-backend")]
1543 cluster_responder: Option<aion_store_haematite::ClusterResponder>,
1544 #[cfg(feature = "haematite-backend")]
1548 cluster_store: Option<Arc<aion_store_haematite::HaematiteStore>>,
1549 #[cfg(feature = "haematite-backend")]
1552 watched_peers: Vec<crate::cluster::WatchedPeer>,
1553 #[cfg(feature = "haematite-backend")]
1557 directory_peers: Vec<crate::routing::DirectoryPeer>,
1558 #[cfg(feature = "haematite-backend")]
1562 self_node_id: Option<String>,
1563}
1564
1565impl ConnectedStore {
1566 fn local(
1573 event_store: Arc<dyn EventStore>,
1574 outbox_store: Option<Arc<dyn OutboxStore>>,
1575 namespace_store: Arc<dyn NamespaceStore>,
1576 ) -> Self {
1577 Self {
1578 event_store,
1579 outbox_store,
1580 namespace_store,
1581 observability_store: None,
1585 bootstrap_coordinator: true,
1586 #[cfg(feature = "haematite-backend")]
1587 cluster_responder: None,
1588 #[cfg(feature = "haematite-backend")]
1589 cluster_store: None,
1590 #[cfg(feature = "haematite-backend")]
1591 watched_peers: Vec::new(),
1592 #[cfg(feature = "haematite-backend")]
1593 directory_peers: Vec::new(),
1594 #[cfg(feature = "haematite-backend")]
1595 self_node_id: None,
1596 }
1597 }
1598}
1599
1600async fn connect_store(config: StoreConfig) -> Result<ConnectedStore, ServerError> {
1610 match config.backend {
1611 StoreBackend::Memory => {
1612 let leaf = Arc::new(aion_store::InMemoryStore::default());
1615 let namespace_store: Arc<dyn NamespaceStore> = leaf.clone();
1616 Ok(ConnectedStore::local(leaf, None, namespace_store))
1617 }
1618 StoreBackend::LibSql => {
1619 #[cfg(feature = "libsql-backend")]
1620 {
1621 connect_libsql_store(config).await
1622 }
1623 #[cfg(not(feature = "libsql-backend"))]
1624 {
1625 let _ = config;
1626 connect_libsql_store_unavailable()
1627 }
1628 }
1629 StoreBackend::Haematite => {
1630 #[cfg(feature = "haematite-backend")]
1631 {
1632 connect_haematite_store(config).await
1633 }
1634 #[cfg(not(feature = "haematite-backend"))]
1635 {
1636 let _ = config;
1637 connect_haematite_store_unavailable()
1638 }
1639 }
1640 }
1641}
1642
1643#[cfg(feature = "libsql-backend")]
1647async fn connect_libsql_store(config: StoreConfig) -> Result<ConnectedStore, ServerError> {
1648 let Some(url) = config.url else {
1649 return Err(ServerError::Config {
1650 message: "store.url must not be empty when store.backend is libsql".to_owned(),
1651 });
1652 };
1653 let store = LibSqlStore::open(url.clone())
1654 .await
1655 .map_err(ServerError::from)?;
1656 store
1657 .validate_event_compatibility()
1658 .await
1659 .map_err(|error| match error {
1660 aion_store::StoreError::Serialization(_) => ServerError::Config {
1661 message: format!(
1662 "Database schema mismatch — delete {url} and restart, or run migrations."
1663 ),
1664 },
1665 other => ServerError::from(other),
1666 })?;
1667 let leaf = Arc::new(store);
1668 let event_store: Arc<dyn EventStore> = leaf.clone();
1669 let namespace_store: Arc<dyn NamespaceStore> = leaf.clone();
1670 let outbox_store: Arc<dyn OutboxStore> = leaf;
1671 Ok(ConnectedStore::local(
1672 event_store,
1673 Some(outbox_store),
1674 namespace_store,
1675 ))
1676}
1677
1678#[cfg(not(feature = "libsql-backend"))]
1682fn connect_libsql_store_unavailable() -> Result<ConnectedStore, ServerError> {
1683 Err(ServerError::Config {
1684 message: "store.backend = libsql requires the aion-server `libsql-backend` feature"
1685 .to_owned(),
1686 })
1687}
1688
1689#[cfg(feature = "haematite-backend")]
1709async fn connect_haematite_store(config: StoreConfig) -> Result<ConnectedStore, ServerError> {
1710 let Some(data_dir) = config.data_dir else {
1711 return Err(ServerError::Config {
1712 message: "store.data_dir must not be empty when store.backend is haematite".to_owned(),
1713 });
1714 };
1715 let shard_count = config.shard_count;
1716 let owned_shards = config.owned_shards.clone();
1717 let cluster = config.cluster.clone();
1718 let watched_peers: Vec<crate::cluster::WatchedPeer> = cluster
1723 .as_ref()
1724 .map(|cluster| {
1725 cluster
1726 .peers
1727 .iter()
1728 .map(|peer| crate::cluster::WatchedPeer {
1729 name: peer.name.clone(),
1730 owned_shards: peer.owned_shards.clone(),
1731 })
1732 .collect()
1733 })
1734 .unwrap_or_default();
1735 let directory_peers: Vec<crate::routing::DirectoryPeer> = cluster
1738 .as_ref()
1739 .map(|cluster| {
1740 cluster
1741 .peers
1742 .iter()
1743 .map(|peer| crate::routing::DirectoryPeer {
1744 name: peer.name.clone(),
1745 owned_shards: peer.owned_shards.clone(),
1746 grpc_addr: peer.grpc_address,
1747 })
1748 .collect()
1749 })
1750 .unwrap_or_default();
1751 let self_node_id: Option<String> = cluster.as_ref().map(|cluster| cluster.node_id.clone());
1754 let (store, responder) =
1758 tokio::task::spawn_blocking(move || build_haematite_store(&data_dir, shard_count, cluster))
1759 .await
1760 .map_err(|error| ServerError::Config {
1761 message: format!("haematite store initialization task failed: {error}"),
1762 })??;
1763
1764 let bootstrap_coordinator = if owned_shards.is_empty() {
1768 true
1769 } else {
1770 store.set_owned_shards(owned_shards.iter().copied());
1771 store.owns_workflow_shard(&aion::schedule_coordinator_workflow_id())
1772 };
1773
1774 let leaf = Arc::new(store);
1775 let event_store: Arc<dyn EventStore> = leaf.clone();
1776 let outbox_store: Arc<dyn OutboxStore> = leaf.clone();
1777 let namespace_store: Arc<dyn NamespaceStore> = leaf.clone();
1781 let observability_store: Arc<dyn aion_store::ObservabilityStore> = leaf.clone();
1786 let cluster_store = responder.as_ref().map(|_| leaf);
1790 let (watched_peers, directory_peers, self_node_id) = if cluster_store.is_some() {
1791 (watched_peers, directory_peers, self_node_id)
1792 } else {
1793 (Vec::new(), Vec::new(), None)
1794 };
1795 Ok(ConnectedStore {
1796 event_store,
1797 outbox_store: Some(outbox_store),
1798 namespace_store,
1799 observability_store: Some(observability_store),
1800 bootstrap_coordinator,
1801 cluster_responder: responder,
1802 cluster_store,
1803 watched_peers,
1804 directory_peers,
1805 self_node_id,
1806 })
1807}
1808
1809#[cfg(feature = "haematite-backend")]
1824fn build_haematite_store(
1825 data_dir: &str,
1826 shard_count: usize,
1827 cluster: Option<crate::config::ClusterConfig>,
1828) -> Result<
1829 (
1830 aion_store_haematite::HaematiteStore,
1831 Option<aion_store_haematite::ClusterResponder>,
1832 ),
1833 ServerError,
1834> {
1835 build_haematite_store_with_hook(data_dir, shard_count, cluster, || Ok(()))
1836}
1837
1838#[cfg(feature = "haematite-backend")]
1839fn build_haematite_store_with_hook(
1840 data_dir: &str,
1841 shard_count: usize,
1842 cluster: Option<crate::config::ClusterConfig>,
1843 before_backend_touch: impl FnOnce() -> Result<(), std::io::Error>,
1844) -> Result<
1845 (
1846 aion_store_haematite::HaematiteStore,
1847 Option<aion_store_haematite::ClusterResponder>,
1848 ),
1849 ServerError,
1850> {
1851 use aion_store_haematite::{ClusterBootstrap, HaematiteStore};
1852
1853 let private_root = crate::filesystem::ConfinedDir::open_or_create(std::path::Path::new(
1861 data_dir,
1862 ))
1863 .map_err(|error| ServerError::Config {
1864 message: format!("unsafe store.data_dir `{data_dir}`: {error}"),
1865 })?;
1866
1867 for shard in 0..shard_count {
1871 private_root
1872 .create_dir_all(std::path::Path::new(&format!("shard-{shard}")))
1873 .map_err(|error| ServerError::Config {
1874 message: format!(
1875 "failed to materialize shard-{shard} under store.data_dir `{data_dir}`: {error}"
1876 ),
1877 })?;
1878 }
1879 private_root
1880 .harden_tree()
1881 .map_err(|error| private_store_mode_error(data_dir, &error))?;
1882
1883 before_backend_touch().map_err(|error| ServerError::Config {
1886 message: format!("store.data_dir pre-open hook failed: {error}"),
1887 })?;
1888
1889 #[cfg(unix)]
1890 let backend_path = private_root
1891 .backend_path()
1892 .map_err(|error| ServerError::Config {
1893 message: format!("failed to resolve held store.data_dir `{data_dir}`: {error}"),
1894 })?;
1895 #[cfg(all(unix, not(any(target_os = "linux", target_os = "android"))))]
1896 crate::filesystem::validate_ambient_backend_ancestors(&backend_path).map_err(|error| {
1897 let (component, reason) = error.into_parts();
1898 ServerError::UnsafeDataRootAncestor {
1899 data_root: backend_path.clone(),
1900 component,
1901 reason,
1902 }
1903 })?;
1904 #[cfg(not(unix))]
1905 let backend_path = std::path::PathBuf::from(data_dir);
1906
1907 let Some(cluster) = cluster else {
1908 let store = if backend_path.join("config.json").exists() {
1909 HaematiteStore::open(&backend_path).map_err(ServerError::from)?
1910 } else {
1911 HaematiteStore::create_with_shard_count(&backend_path, shard_count)
1912 .map_err(ServerError::from)?
1913 };
1914 store.materialize_all_shards().map_err(ServerError::from)?;
1915 private_root
1916 .harden_tree()
1917 .map_err(|error| private_store_mode_error(data_dir, &error))?;
1918 let store = store.retain_data_root_capability(private_root);
1919 return Ok((store, None));
1920 };
1921
1922 let boot = ClusterBootstrap {
1923 node_id: cluster.node_id,
1924 bind_address: cluster.bind_address,
1925 members: cluster.members,
1926 peers: cluster
1927 .peers
1928 .into_iter()
1929 .map(|peer| (peer.name, peer.address))
1930 .collect(),
1931 timeout: HAEMATITE_CLUSTER_OP_TIMEOUT,
1932 };
1933 let (store, responder) =
1934 HaematiteStore::open_or_create_distributed(&backend_path, shard_count, boot)
1935 .map_err(ServerError::from)?;
1936 store.materialize_all_shards().map_err(ServerError::from)?;
1937 private_root
1938 .harden_tree()
1939 .map_err(|error| private_store_mode_error(data_dir, &error))?;
1940 let store = store.retain_data_root_capability(private_root);
1941 Ok((store, Some(responder)))
1942}
1943
1944#[cfg(feature = "haematite-backend")]
1945fn private_store_mode_error(data_dir: &str, error: &std::io::Error) -> ServerError {
1946 ServerError::Config {
1947 message: format!(
1948 "failed to apply private modes under store.data_dir `{data_dir}`: {error}"
1949 ),
1950 }
1951}
1952
1953#[cfg(feature = "haematite-backend")]
1955const HAEMATITE_CLUSTER_OP_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(5);
1956
1957#[cfg(not(feature = "haematite-backend"))]
1961fn connect_haematite_store_unavailable() -> Result<ConnectedStore, ServerError> {
1962 Err(ServerError::Config {
1963 message: "store.backend = haematite requires the aion-server `haematite-backend` feature"
1964 .to_owned(),
1965 })
1966}
1967
1968#[cfg(not(feature = "liminal-transport"))]
1977#[derive(Clone, Debug)]
1978struct NullInterventionTransport;
1979
1980#[cfg(not(feature = "liminal-transport"))]
1981#[async_trait::async_trait]
1982impl crate::worker::InterventionTransport for NullInterventionTransport {
1983 async fn push(
1984 &self,
1985 _worker: &crate::worker::WorkerHandle,
1986 _command: aion_core::InterventionCommand,
1987 ) -> Result<aion_core::InterventionOutcome, ServerError> {
1988 Err(ServerError::worker_connection_lost(
1989 "intervention",
1990 "no intervention push transport is compiled in".to_owned(),
1991 ))
1992 }
1993}
1994
1995#[cfg(test)]
1996mod tests {
1997 use std::{net::SocketAddr, time::Duration};
1998
1999 use aion_store::InMemoryStore;
2000
2001 use super::ServerState;
2002 use crate::config::{
2003 AuthConfig, AuthoringConfig, DeployConfig, DevConfig, ListenConfig, MetricsConfig,
2004 NamespaceConfig, NamespaceMode, OpsConsoleAssetSource, OpsConsoleConfig, OutboxConfig,
2005 RuntimeConfig, WebSocketConfig, WorkerConfig,
2006 };
2007
2008 fn runtime_config() -> RuntimeConfig {
2009 RuntimeConfig {
2010 listen: ListenConfig {
2011 grpc: SocketAddr::from(([127, 0, 0, 1], 50051)),
2012 http: SocketAddr::from(([127, 0, 0, 1], 8080)),
2013 },
2014 tls: None,
2015 auth: AuthConfig {
2016 enabled: false,
2017 jwks_url: None,
2018 jwks_refresh_seconds: 300,
2019 },
2020 ops_console: OpsConsoleConfig {
2021 source: OpsConsoleAssetSource::Embedded,
2022 },
2023 namespace: NamespaceConfig {
2024 mode: NamespaceMode::SharedEngine,
2025 },
2026 worker: WorkerConfig {
2027 heartbeat_window: Duration::from_secs(30),
2028 ..WorkerConfig::default()
2029 },
2030 websocket: WebSocketConfig {
2031 outbound_buffer_bound: 32,
2032 event_broadcast_capacity: Some(64),
2033 cluster_broadcast_capacity: Some(64),
2034 },
2035 workflow_packages: Vec::new(),
2036 deploy: DeployConfig::default(),
2037 authoring: AuthoringConfig::default(),
2038 dev: DevConfig::default(),
2039 outbox: OutboxConfig::default(),
2040 observability: crate::config::ObservabilityConfig::default(),
2041 scheduler_threads: 1,
2042 query_timeout: Some(Duration::from_secs(10)),
2043 default_namespace: "default".to_owned(),
2044 auto_create: crate::config::AutoCreate::Open,
2045 max_in_flight_activities: crate::config::DEFAULT_MAX_IN_FLIGHT_ACTIVITIES,
2046 drain_timeout: Duration::from_secs(30),
2047 metrics: MetricsConfig { enabled: true },
2048 owned_shards: Vec::new(),
2049 cors_allowed_origins: Vec::new(),
2050 }
2051 }
2052
2053 #[tokio::test]
2054 async fn builds_state_with_in_memory_store() -> Result<(), Box<dyn std::error::Error>> {
2055 let state =
2056 ServerState::build_with_store(InMemoryStore::default(), runtime_config()).await?;
2057
2058 std::hint::black_box(state.namespace_guard());
2059 std::hint::black_box(state.worker_registry());
2060
2061 Ok(())
2062 }
2063
2064 #[tokio::test]
2072 async fn unserved_queues_surfaces_a_parked_dispatch_and_clears_it()
2073 -> Result<(), Box<dyn std::error::Error>> {
2074 use aion::{ActivityDispatch, ActivityDispatcher as _};
2075 use aion_core::{ActivityId, RunId, WorkflowId};
2076 use std::sync::Arc;
2077
2078 let state =
2079 ServerState::build_with_store(InMemoryStore::default(), runtime_config()).await?;
2080 assert!(
2081 state.unserved_queues()?.is_empty(),
2082 "a calm boot has no unserved queues"
2083 );
2084 assert!(state.queue_declarations().is_installed());
2088 assert_eq!(
2089 state
2090 .queue_declarations()
2091 .declaration_for("nobody-serves-this"),
2092 crate::worker::QueueDeclaration::Unknown
2093 );
2094
2095 let dispatcher = Arc::new(
2096 crate::worker::WorkerActivityDispatcher::new(
2097 state.worker_registry().clone(),
2098 "default",
2099 crate::worker::HeartbeatTracker::new(Duration::from_secs(5)),
2100 )
2101 .with_queue_state(state.queue_service_state().clone())
2102 .with_queue_declarations(state.queue_declarations().clone()),
2103 );
2104 let workflow_id = WorkflowId::new_v4();
2105 let request = ActivityDispatch {
2106 namespace: "default".to_owned(),
2107 task_queue: "nobody-serves-this".to_owned(),
2108 node: None,
2109 workflow_id: workflow_id.clone(),
2110 run_id: RunId::new_v4(),
2111 activity_id: ActivityId::from_sequence_position(0),
2112 name: "greet".to_owned(),
2113 input: "{}".to_owned(),
2114 config: "{}".to_owned(),
2115 attempt: 1,
2116 advisory: false,
2117 labels: std::collections::BTreeMap::new(),
2118 };
2119 let parked = std::thread::spawn(move || dispatcher.dispatch(request));
2120
2121 let mut unserved = Vec::new();
2122 for _ in 0..30 {
2123 unserved = state.unserved_queues()?;
2124 if !unserved.is_empty() {
2125 break;
2126 }
2127 tokio::time::sleep(Duration::from_millis(100)).await;
2128 }
2129 assert_eq!(unserved.len(), 1, "the parked dispatch is not surfaced");
2130 assert_eq!(
2131 unserved[0].reason,
2132 crate::worker::QueueServiceReason::NoLivePollers,
2133 "an empty catalog must not be read as a structural refusal"
2134 );
2135 assert_eq!(unserved[0].key.task_queue, "nobody-serves-this");
2136 assert_eq!(unserved[0].waiting.len(), 1);
2137 assert_eq!(unserved[0].waiting[0].workflow_id, workflow_id);
2138
2139 let (worker_tx, worker_rx) = tokio::sync::mpsc::channel(1);
2141 drop(worker_rx);
2142 let registration = state.worker_registry().register_namespaces(
2143 [String::from("default")],
2144 "nobody-serves-this",
2145 None,
2146 [String::from("greet")].iter(),
2147 worker_tx,
2148 )?;
2149 let outcome = parked.join().map_err(|_| "parked dispatch panicked")?;
2150 assert!(outcome.is_err(), "the released dispatch must resolve");
2151 assert!(
2152 state.unserved_queues()?.is_empty(),
2153 "a resolved dispatch must leave the unserved state"
2154 );
2155 registration.deregister()?;
2156 Ok(())
2157 }
2158
2159 #[tokio::test]
2160 async fn namespace_store_is_reachable_and_functional_after_default_boot()
2161 -> Result<(), Box<dyn std::error::Error>> {
2162 use aion_store::{MintOutcome, NamespaceOrigin};
2163
2164 let state =
2168 ServerState::build_with_store(InMemoryStore::default(), runtime_config()).await?;
2169
2170 let store = state.namespace_store();
2171
2172 let outcome = store
2174 .register_namespace("orders", NamespaceOrigin::WorkerMint)
2175 .await?;
2176 assert_eq!(
2177 outcome,
2178 MintOutcome::Created,
2179 "the first reference to a namespace mints it"
2180 );
2181
2182 let again = store
2184 .register_namespace("orders", NamespaceOrigin::WorkerMint)
2185 .await?;
2186 assert_eq!(
2187 again,
2188 MintOutcome::AlreadyExisted,
2189 "a second reference touches the existing record rather than re-creating it"
2190 );
2191
2192 let fetched = store.get_namespace("orders").await?;
2194 let record = fetched.ok_or("registered namespace must be retrievable via get_namespace")?;
2195 assert_eq!(record.name, "orders");
2196 assert_eq!(record.origin, NamespaceOrigin::WorkerMint);
2197
2198 let listed = store.list_namespaces().await?;
2200 assert!(
2201 listed.iter().any(|record| record.name == "orders"),
2202 "list_namespaces returns the minted namespace"
2203 );
2204
2205 Ok(())
2206 }
2207
2208 #[cfg(feature = "haematite-backend")]
2209 #[tokio::test(flavor = "multi_thread")]
2210 async fn connect_store_haematite_round_trips_through_event_store()
2211 -> Result<(), Box<dyn std::error::Error>> {
2212 use aion_core::{ContentType, EventEnvelope, PackageVersion, Payload, RunId, WorkflowId};
2213 use aion_store::WriteToken;
2214 use chrono::Utc;
2215
2216 use crate::config::{StoreBackend, StoreConfig};
2217
2218 let data_dir = crate::test_support::private_tempdir()?;
2219 let connected = super::connect_store(StoreConfig {
2223 backend: StoreBackend::Haematite,
2224 url: None,
2225 owned_shards: Vec::new(),
2226 data_dir: Some(data_dir.path().to_string_lossy().into_owned()),
2227 shard_count: 1,
2228 cluster: None,
2229 })
2230 .await?;
2231 let event_store = connected.event_store;
2232 assert!(
2233 connected.outbox_store.is_some(),
2234 "the haematite backend shares its leaf store as the dispatcher's outbox store"
2235 );
2236 assert!(
2237 connected.bootstrap_coordinator,
2238 "a single-node haematite boot owns all shards and bootstraps the coordinator"
2239 );
2240 assert!(
2241 connected.cluster_responder.is_none(),
2242 "a single-node (no [cluster]) haematite boot has no distributed responder"
2243 );
2244
2245 let workflow_id = WorkflowId::new_v4();
2246 let event = aion_core::Event::WorkflowStarted {
2247 envelope: EventEnvelope {
2248 seq: 1,
2249 recorded_at: Utc::now(),
2250 workflow_id: workflow_id.clone(),
2251 },
2252 workflow_type: String::from("checkout"),
2253 input: Payload::new(ContentType::Json, b"{}".to_vec()),
2254 run_id: RunId::new_v4(),
2255 parent_run_id: None,
2256 package_version: PackageVersion::new("a".repeat(64)),
2257 };
2258 event_store
2259 .append(
2260 WriteToken::recorder(),
2261 &workflow_id,
2262 std::slice::from_ref(&event),
2263 0,
2264 )
2265 .await?;
2266 let history = event_store.read_history(&workflow_id).await?;
2267 assert_eq!(
2268 history.len(),
2269 1,
2270 "an event appended through the server's dyn EventStore reads back"
2271 );
2272 Ok(())
2273 }
2274
2275 #[cfg(all(feature = "haematite-backend", unix))]
2276 #[test]
2277 fn haematite_root_swap_before_first_backend_touch_cannot_redirect_writes()
2278 -> Result<(), Box<dyn std::error::Error>> {
2279 use std::os::unix::fs::symlink;
2280
2281 let sandbox = crate::test_support::private_tempdir()?;
2282 let configured_root = sandbox.path().join("data");
2283 let held_root = sandbox.path().join("held-data");
2284 let outside = sandbox.path().join("outside");
2285 std::fs::create_dir(&outside)?;
2286 let configured = configured_root
2287 .to_str()
2288 .ok_or("temporary data path was not UTF-8")?;
2289
2290 let (store, responder) =
2291 super::build_haematite_store_with_hook(configured, 4, None, || {
2292 std::fs::rename(&configured_root, &held_root)?;
2297 symlink(&outside, &configured_root)?;
2298 Ok(())
2299 })?;
2300 assert!(responder.is_none());
2301
2302 let outside_entries = std::fs::read_dir(&outside)?.collect::<Result<Vec<_>, _>>()?;
2303 assert!(
2304 outside_entries.is_empty(),
2305 "Haematite followed the replaced ambient root and wrote outside"
2306 );
2307 assert!(held_root.join("config.json").is_file());
2308 for shard in 0..4 {
2309 let shard_path = held_root.join(format!("shard-{shard}"));
2310 assert!(shard_path.is_dir(), "shard {shard} was not materialized");
2311 assert!(
2312 std::fs::read_dir(&shard_path)?
2313 .next()
2314 .transpose()?
2315 .is_some(),
2316 "shard {shard} did not run Haematite's materialization path"
2317 );
2318 }
2319
2320 drop(store);
2321 Ok(())
2322 }
2323
2324 #[cfg(all(
2325 feature = "haematite-backend",
2326 any(target_os = "linux", target_os = "android")
2327 ))]
2328 #[tokio::test]
2329 async fn proc_fd_backend_path_survives_a_post_startup_root_swap()
2330 -> Result<(), Box<dyn std::error::Error>> {
2331 use std::os::unix::fs::symlink;
2332
2333 use aion_core::{ContentType, EventEnvelope, PackageVersion, Payload, RunId, WorkflowId};
2334 use aion_store::{WritableEventStore as _, WriteToken};
2335 use chrono::Utc;
2336
2337 let sandbox = crate::test_support::private_tempdir()?;
2338 let configured_root = sandbox.path().join("data");
2339 let held_root = sandbox.path().join("held-data");
2340 let capture = sandbox.path().join("capture");
2341 std::fs::create_dir(&capture)?;
2342 let configured = configured_root
2343 .to_str()
2344 .ok_or("temporary data path was not UTF-8")?;
2345
2346 let (store, responder) = super::build_haematite_store(configured, 4, None)?;
2347 assert!(responder.is_none());
2348 std::fs::rename(&configured_root, &held_root)?;
2349 symlink(&capture, &configured_root)?;
2350
2351 let workflow_id = WorkflowId::new_v4();
2352 let event = aion_core::Event::WorkflowStarted {
2353 envelope: EventEnvelope {
2354 seq: 1,
2355 recorded_at: Utc::now(),
2356 workflow_id: workflow_id.clone(),
2357 },
2358 workflow_type: String::from("post-startup-root-swap"),
2359 input: Payload::new(ContentType::Json, b"{}".to_vec()),
2360 run_id: RunId::new_v4(),
2361 parent_run_id: None,
2362 package_version: PackageVersion::new("a".repeat(64)),
2363 };
2364 store
2365 .append(
2366 WriteToken::recorder(),
2367 &workflow_id,
2368 std::slice::from_ref(&event),
2369 0,
2370 )
2371 .await?;
2372
2373 let captured = std::fs::read_dir(&capture)?.collect::<Result<Vec<_>, _>>()?;
2374 assert!(
2375 captured.is_empty(),
2376 "post-startup append followed the replacement symlink into capture"
2377 );
2378 assert!(held_root.join("config.json").is_file());
2379 drop(store);
2380 Ok(())
2381 }
2382
2383 #[cfg(all(
2384 feature = "haematite-backend",
2385 unix,
2386 not(any(target_os = "linux", target_os = "android"))
2387 ))]
2388 #[test]
2389 fn path_ambient_haematite_refuses_group_or_world_writable_ancestors()
2390 -> Result<(), Box<dyn std::error::Error>> {
2391 use std::os::unix::fs::PermissionsExt as _;
2392
2393 let sandbox = crate::test_support::private_tempdir()?;
2394 std::fs::set_permissions(sandbox.path(), std::fs::Permissions::from_mode(0o700))?;
2395
2396 for mode in [0o770, 0o1777] {
2397 let shared = sandbox.path().join(format!("shared-{mode:o}"));
2398 let data_root = shared.join("data");
2399 std::fs::create_dir(&shared)?;
2400 std::fs::set_permissions(&shared, std::fs::Permissions::from_mode(mode))?;
2401 std::fs::create_dir(&data_root)?;
2402 std::fs::set_permissions(&data_root, std::fs::Permissions::from_mode(0o700))?;
2403 let configured = data_root
2404 .to_str()
2405 .ok_or("temporary data path was not UTF-8")?;
2406
2407 let Err(error) = super::build_haematite_store(configured, 4, None) else {
2408 return Err(format!("mode {mode:04o} ancestor was accepted").into());
2409 };
2410 let message = error.to_string();
2411 let crate::ServerError::UnsafeDataRootAncestor {
2412 data_root: resolved_root,
2413 component,
2414 reason,
2415 } = error
2416 else {
2417 return Err(format!("expected typed unsafe-ancestor error, got {message}").into());
2418 };
2419 assert_eq!(resolved_root, std::fs::canonicalize(&data_root)?);
2420 assert_eq!(component, std::fs::canonicalize(&shared)?);
2421 assert!(
2422 reason.contains(&format!("mode {mode:04o}")),
2423 "unexpected reason: {reason}"
2424 );
2425 if mode & 0o1000 != 0 {
2426 assert!(reason.contains("sticky bit is not accepted"));
2427 }
2428 assert!(message.contains("private Aion home"));
2429 assert!(
2430 !data_root.join("config.json").exists(),
2431 "Haematite touched its ambient path before the refusal"
2432 );
2433 }
2434 Ok(())
2435 }
2436
2437 #[cfg(all(feature = "haematite-backend", target_os = "macos"))]
2438 #[test]
2439 fn path_ambient_haematite_refuses_mutating_allow_acl_ancestor()
2440 -> Result<(), Box<dyn std::error::Error>> {
2441 use std::os::unix::fs::PermissionsExt as _;
2442
2443 let sandbox = crate::test_support::private_tempdir()?;
2444 std::fs::set_permissions(sandbox.path(), std::fs::Permissions::from_mode(0o700))?;
2445 let shared = sandbox.path().join("acl-shared");
2446 let data_root = shared.join("data");
2447 std::fs::create_dir(&shared)?;
2448 std::fs::set_permissions(&shared, std::fs::Permissions::from_mode(0o700))?;
2449 let acl = "everyone allow list,search,add_file,add_subdirectory,delete_child";
2450 let status = std::process::Command::new("chmod")
2451 .arg("+a")
2452 .arg(acl)
2453 .arg(&shared)
2454 .status()?;
2455 assert!(status.success(), "failed to install Darwin regression ACL");
2456 let configured = data_root
2457 .to_str()
2458 .ok_or("temporary data path was not UTF-8")?;
2459
2460 let result = super::build_haematite_store(configured, 4, None);
2461 let cleanup = std::process::Command::new("chmod")
2462 .arg("-RN")
2463 .arg(&shared)
2464 .status()?;
2465 assert!(cleanup.success(), "failed to clean Darwin regression ACL");
2466
2467 let Err(error) = result else {
2468 return Err("mutating non-euid allow ACL ancestor was accepted".into());
2469 };
2470 let message = error.to_string();
2471 let crate::ServerError::UnsafeDataRootAncestor {
2472 component, reason, ..
2473 } = error
2474 else {
2475 return Err(format!("expected typed unsafe-ancestor error, got {message}").into());
2476 };
2477 assert_eq!(component, std::fs::canonicalize(&shared)?);
2478 assert!(
2479 reason.contains("allow"),
2480 "reason did not name the ACE: {reason}"
2481 );
2482 assert!(
2483 reason.contains("everyone"),
2484 "reason did not name the ACE principal: {reason}"
2485 );
2486 assert!(
2487 !data_root.join("config.json").exists(),
2488 "Haematite touched its ambient path before the ACL refusal"
2489 );
2490 Ok(())
2491 }
2492
2493 #[cfg(all(feature = "haematite-backend", target_os = "macos"))]
2494 #[test]
2495 fn path_ambient_haematite_accepts_the_euid_uuid_allow_ace()
2496 -> Result<(), Box<dyn std::error::Error>> {
2497 use std::os::unix::fs::PermissionsExt as _;
2498
2499 use exacl::{AclEntry, AclOption, Perm};
2500
2501 let sandbox = crate::test_support::private_tempdir()?;
2502 std::fs::set_permissions(sandbox.path(), std::fs::Permissions::from_mode(0o700))?;
2503 let private_parent = sandbox.path().join("euid-uuid-allow");
2504 let data_root = private_parent.join("data");
2505 std::fs::create_dir(&private_parent)?;
2506 std::fs::set_permissions(&private_parent, std::fs::Permissions::from_mode(0o700))?;
2507
2508 let server_uid = rustix::process::geteuid().as_raw();
2509 let ace_qualifier = crate::filesystem::darwin_user_uuid_for_test(server_uid)?;
2510 let entry = AclEntry::allow_user(
2511 &ace_qualifier.to_string(),
2512 Perm::EXECUTE | Perm::WRITE | Perm::APPEND | Perm::DELETE_CHILD,
2513 None,
2514 );
2515 exacl::setfacl(
2516 &[private_parent.as_path()],
2517 &[entry],
2518 AclOption::SYMLINK_ACL,
2519 )?;
2520 let configured = data_root
2521 .to_str()
2522 .ok_or("temporary data path was not UTF-8")?;
2523
2524 let result = super::build_haematite_store(configured, 4, None);
2525 let cleanup = std::process::Command::new("chmod")
2526 .arg("-RN")
2527 .arg(&private_parent)
2528 .status()?;
2529 assert!(cleanup.success(), "failed to clean euid UUID allow ACL");
2530
2531 let (store, responder) = result?;
2532 assert!(responder.is_none());
2533 assert!(data_root.join("config.json").is_file());
2534 drop(store);
2535 Ok(())
2536 }
2537
2538 #[cfg(all(feature = "haematite-backend", target_os = "macos"))]
2539 #[test]
2540 fn path_ambient_haematite_refuses_a_non_euid_user_uuid_allow_ace()
2541 -> Result<(), Box<dyn std::error::Error>> {
2542 use std::os::unix::fs::PermissionsExt as _;
2543
2544 use exacl::{AclEntry, AclOption, Perm};
2545
2546 let sandbox = crate::test_support::private_tempdir()?;
2547 std::fs::set_permissions(sandbox.path(), std::fs::Permissions::from_mode(0o700))?;
2548 let shared = sandbox.path().join("non-euid-uuid-allow");
2549 let data_root = shared.join("data");
2550 std::fs::create_dir(&shared)?;
2551 std::fs::set_permissions(&shared, std::fs::Permissions::from_mode(0o700))?;
2552
2553 let server_uid = rustix::process::geteuid().as_raw();
2554 let foreign_uid = u32::from(server_uid == 0);
2555 let foreign_qualifier = crate::filesystem::darwin_user_uuid_for_test(foreign_uid)?;
2556 let entry = AclEntry::allow_user(
2557 &foreign_qualifier.to_string(),
2558 Perm::EXECUTE | Perm::WRITE | Perm::APPEND | Perm::DELETE_CHILD,
2559 None,
2560 );
2561 exacl::setfacl(&[shared.as_path()], &[entry], AclOption::SYMLINK_ACL)?;
2562 let configured = data_root
2563 .to_str()
2564 .ok_or("temporary data path was not UTF-8")?;
2565
2566 let result = super::build_haematite_store(configured, 4, None);
2567 let cleanup = std::process::Command::new("chmod")
2568 .arg("-RN")
2569 .arg(&shared)
2570 .status()?;
2571 assert!(cleanup.success(), "failed to clean non-euid UUID allow ACL");
2572
2573 let Err(error) = result else {
2574 return Err("mutating non-euid user UUID allow ACE was accepted".into());
2575 };
2576 let message = error.to_string();
2577 let crate::ServerError::UnsafeDataRootAncestor {
2578 component, reason, ..
2579 } = error
2580 else {
2581 return Err(format!("expected typed unsafe-ancestor error, got {message}").into());
2582 };
2583 assert_eq!(component, std::fs::canonicalize(&shared)?);
2584 assert!(
2585 reason.contains("allow") && reason.contains(&format!("server euid {server_uid}")),
2586 "reason did not name the rejected ACE: {reason}"
2587 );
2588 assert!(
2589 !data_root.join("config.json").exists(),
2590 "Haematite touched its ambient path before the UUID ACL refusal"
2591 );
2592 Ok(())
2593 }
2594
2595 #[cfg(all(feature = "haematite-backend", target_os = "macos"))]
2596 #[test]
2597 fn path_ambient_haematite_accepts_a_deny_only_acl_ancestor()
2598 -> Result<(), Box<dyn std::error::Error>> {
2599 use std::os::unix::fs::PermissionsExt as _;
2600
2601 let sandbox = crate::test_support::private_tempdir()?;
2602 std::fs::set_permissions(sandbox.path(), std::fs::Permissions::from_mode(0o700))?;
2603 let private_parent = sandbox.path().join("deny-only");
2604 let data_root = private_parent.join("data");
2605 std::fs::create_dir(&private_parent)?;
2606 std::fs::set_permissions(&private_parent, std::fs::Permissions::from_mode(0o700))?;
2607 let status = std::process::Command::new("chmod")
2608 .arg("+a")
2609 .arg("everyone deny delete")
2610 .arg(&private_parent)
2611 .status()?;
2612 assert!(status.success(), "failed to install Darwin deny-only ACL");
2613 let configured = data_root
2614 .to_str()
2615 .ok_or("temporary data path was not UTF-8")?;
2616
2617 let result = super::build_haematite_store(configured, 4, None);
2618 let cleanup = std::process::Command::new("chmod")
2619 .arg("-RN")
2620 .arg(&private_parent)
2621 .status()?;
2622 assert!(cleanup.success(), "failed to clean Darwin deny-only ACL");
2623
2624 let (store, responder) = result?;
2625 assert!(responder.is_none());
2626 assert!(data_root.join("config.json").is_file());
2627 drop(store);
2628 Ok(())
2629 }
2630
2631 #[cfg(all(feature = "haematite-backend", target_os = "macos"))]
2632 #[test]
2633 fn path_ambient_haematite_accepts_the_stock_home_acl_chain()
2634 -> Result<(), Box<dyn std::error::Error>> {
2635 use std::os::unix::fs::PermissionsExt as _;
2636 use users::os::unix::UserExt as _;
2637
2638 let effective_uid = rustix::process::geteuid().as_raw();
2639 let effective_user = users::get_user_by_uid(effective_uid)
2640 .ok_or_else(|| format!("server euid {effective_uid} has no account record"))?;
2641 let sandbox = tempfile::Builder::new()
2642 .prefix(".aion-acl-home-proof-")
2643 .tempdir_in(effective_user.home_dir())?;
2644 std::fs::set_permissions(sandbox.path(), std::fs::Permissions::from_mode(0o700))?;
2645 let data_root = sandbox.path().join("data");
2646 let configured = data_root
2647 .to_str()
2648 .ok_or("temporary data path was not UTF-8")?;
2649
2650 let (store, responder) = super::build_haematite_store(configured, 4, None)?;
2651 assert!(responder.is_none());
2652 assert!(data_root.join("config.json").is_file());
2653 drop(store);
2654 Ok(())
2655 }
2656
2657 #[cfg(all(
2658 feature = "haematite-backend",
2659 unix,
2660 not(any(target_os = "linux", target_os = "android"))
2661 ))]
2662 #[test]
2663 fn path_ambient_haematite_accepts_an_owner_controlled_chain()
2664 -> Result<(), Box<dyn std::error::Error>> {
2665 use std::os::unix::fs::PermissionsExt as _;
2666
2667 let sandbox = crate::test_support::private_tempdir()?;
2668 std::fs::set_permissions(sandbox.path(), std::fs::Permissions::from_mode(0o700))?;
2669 let private_parent = sandbox.path().join("private");
2670 let data_root = private_parent.join("data");
2671 std::fs::create_dir(&private_parent)?;
2672 std::fs::set_permissions(&private_parent, std::fs::Permissions::from_mode(0o700))?;
2673 let configured = data_root
2674 .to_str()
2675 .ok_or("temporary data path was not UTF-8")?;
2676
2677 let (store, responder) = super::build_haematite_store(configured, 4, None)?;
2678 assert!(responder.is_none());
2679 assert!(data_root.join("config.json").is_file());
2680 for shard in 0..4 {
2681 assert!(data_root.join(format!("shard-{shard}")).is_dir());
2682 }
2683 drop(store);
2684 Ok(())
2685 }
2686
2687 #[tokio::test]
2688 async fn connect_store_memory_backend_exposes_no_outbox_store()
2689 -> Result<(), Box<dyn std::error::Error>> {
2690 use crate::config::{StoreBackend, StoreConfig};
2691
2692 let connected = super::connect_store(StoreConfig {
2695 backend: StoreBackend::Memory,
2696 url: None,
2697 owned_shards: Vec::new(),
2698 data_dir: None,
2699 shard_count: 1,
2700 cluster: None,
2701 })
2702 .await?;
2703 assert!(
2704 connected.outbox_store.is_none(),
2705 "the in-memory backend exposes no outbox store"
2706 );
2707 Ok(())
2708 }
2709
2710 #[cfg(feature = "libsql-backend")]
2714 #[tokio::test]
2715 async fn connect_store_shares_outbox_store_only_for_libsql()
2716 -> Result<(), Box<dyn std::error::Error>> {
2717 use crate::config::{StoreBackend, StoreConfig};
2718
2719 let path = std::env::temp_dir().join(format!(
2724 "aion-connect-store-{}-{}.db",
2725 std::process::id(),
2726 std::time::SystemTime::now()
2727 .duration_since(std::time::UNIX_EPOCH)
2728 .map(|elapsed| elapsed.as_nanos())
2729 .unwrap_or_default()
2730 ));
2731 let connected = super::connect_store(StoreConfig {
2732 backend: StoreBackend::LibSql,
2733 url: Some(path.to_string_lossy().into_owned()),
2734 owned_shards: Vec::new(),
2735 data_dir: None,
2736 shard_count: 1,
2737 cluster: None,
2738 })
2739 .await?;
2740 assert!(
2741 connected.outbox_store.is_some(),
2742 "the libSQL backend shares its leaf store as the dispatcher's outbox store"
2743 );
2744 Ok(())
2745 }
2746
2747 #[tokio::test]
2748 async fn state_build_fails_without_event_broadcast_capacity()
2749 -> Result<(), Box<dyn std::error::Error>> {
2750 let mut runtime = runtime_config();
2751 runtime.websocket.event_broadcast_capacity = None;
2752
2753 let error = ServerState::build_with_store(InMemoryStore::default(), runtime)
2754 .await
2755 .err()
2756 .ok_or("state build must fail when event streaming is unsized")?;
2757
2758 assert!(error.is_config(), "expected a config error, got {error}");
2759 assert!(
2760 error
2761 .to_string()
2762 .contains("websocket.event_broadcast_capacity"),
2763 "error must name the missing key: {error}"
2764 );
2765 Ok(())
2766 }
2767
2768 #[tokio::test]
2769 async fn state_build_fails_without_query_timeout() -> Result<(), Box<dyn std::error::Error>> {
2770 let mut runtime = runtime_config();
2771 runtime.query_timeout = None;
2772
2773 let error = ServerState::build_with_store(InMemoryStore::default(), runtime)
2774 .await
2775 .err()
2776 .ok_or("state build must fail when the query reply deadline is unset")?;
2777
2778 assert!(error.is_config(), "expected a config error, got {error}");
2779 assert!(
2780 error.to_string().contains("runtime.query_timeout_ms"),
2781 "error must name the missing key: {error}"
2782 );
2783 assert!(
2784 error.to_string().contains("AION_RUNTIME_QUERY_TIMEOUT_MS"),
2785 "error must name the environment override: {error}"
2786 );
2787 Ok(())
2788 }
2789
2790 #[tokio::test]
2791 async fn state_build_fails_with_zero_query_timeout() -> Result<(), Box<dyn std::error::Error>> {
2792 let mut runtime = runtime_config();
2793 runtime.query_timeout = Some(Duration::ZERO);
2794
2795 let error = ServerState::build_with_store(InMemoryStore::default(), runtime)
2796 .await
2797 .err()
2798 .ok_or("state build must fail when the query reply deadline is zero")?;
2799
2800 assert!(error.is_config(), "expected a config error, got {error}");
2801 assert!(
2802 error.to_string().contains("runtime.query_timeout_ms"),
2803 "error must name the zero-valued key: {error}"
2804 );
2805 Ok(())
2806 }
2807}