1use std::{path::PathBuf, sync::Arc};
4
5use aion::{
6 ActivityDispatcher, EngineBuilder, RuntimeHandle, SignalRouter, signal::ConcreteSignalRouter,
7};
8use aion_store::{EventStore, NamespaceStore, OutboxStore, WorkerDeploymentStore};
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 supervisor::WorkerSupervisor,
28 },
29};
30
31fn new_supervisor(
42 store: &Arc<dyn WorkerDeploymentStore>,
43 publisher: &crate::cluster_publisher::ClusterEventPublisher,
44) -> Arc<WorkerSupervisor> {
45 Arc::new(WorkerSupervisor::new(Arc::clone(store), publisher.clone()))
46}
47
48#[derive(Clone)]
50pub struct ServerState {
51 inner: Arc<ServerStateInner>,
52}
53
54struct ServerStateInner {
55 namespace_guard: NamespaceGuard,
56 runtime: RuntimeConfig,
57 worker_registry: ConnectedWorkerRegistry,
58 pending_activities: PendingActivities,
59 heartbeat_tracker: HeartbeatTracker,
60 drain_state: DrainState,
61 metrics: Option<Metrics>,
62 health: Option<HealthState>,
63 activity_mock_registry: Option<ActivityMockRegistry>,
66 outbox_store: Option<Arc<dyn OutboxStore>>,
70 namespace_store: Arc<dyn NamespaceStore>,
79 worker_deployment_store: Arc<dyn WorkerDeploymentStore>,
81 worker_supervisor: Arc<WorkerSupervisor>,
86 outbox_wake: Arc<tokio::sync::Notify>,
93 cluster_publisher: crate::cluster_publisher::ClusterEventPublisher,
97 transcript_publisher: crate::activity_publisher::ActivityEventPublisher,
106 attempt_owners: crate::worker::AttemptOwnerIndex,
113 queue_service_state: crate::worker::QueueServiceState,
119 queue_declarations: crate::worker::QueueDeclarationSource,
123 workspace_root: crate::worker::WorkspaceRoot,
131 cluster_self_node: Option<String>,
136 #[cfg(feature = "haematite-backend")]
141 cluster_responder: Option<aion_store_haematite::ClusterResponder>,
142 #[cfg(feature = "haematite-backend")]
145 cluster_store: Option<Arc<aion_store_haematite::HaematiteStore>>,
146 #[cfg(feature = "haematite-backend")]
149 watched_peers: Vec<crate::cluster::WatchedPeer>,
150 #[cfg(feature = "haematite-backend")]
155 shard_directory: Option<Arc<crate::routing::StaticShardDirectory>>,
156 #[cfg(feature = "haematite-backend")]
160 request_forwarder: Option<Arc<dyn crate::routing::RequestForwarder>>,
161 #[cfg(feature = "auth")]
162 jwks_cache: Option<JwksCache>,
163}
164
165impl ServerState {
166 const FALLBACK_CLUSTER_BROADCAST_CAPACITY: std::num::NonZeroUsize =
175 match std::num::NonZeroUsize::new(64) {
176 Some(value) => value,
177 None => std::num::NonZeroUsize::MIN,
178 };
179
180 pub async fn build(config: ServerConfig) -> Result<Self, ServerError> {
187 let (store_config, runtime) = config.into_parts();
188 let connected = connect_store(store_config).await?;
189 Self::build_with_connected_store(connected, runtime).await
190 }
191
192 pub async fn build_with_store<S>(store: S, runtime: RuntimeConfig) -> Result<Self, ServerError>
198 where
199 S: EventStore + NamespaceStore + WorkerDeploymentStore,
200 {
201 let leaf = Arc::new(store);
205 let namespace_store: Arc<dyn NamespaceStore> = leaf.clone();
206 let worker_deployment_store: Arc<dyn WorkerDeploymentStore> = leaf.clone();
207 Self::build_with_connected_store(
208 ConnectedStore::local(leaf, None, namespace_store, worker_deployment_store),
209 runtime,
210 )
211 .await
212 }
213
214 async fn build_with_connected_store(
215 connected: ConnectedStore,
216 runtime: RuntimeConfig,
217 ) -> Result<Self, ServerError> {
218 let cluster_self_node = connected.cluster_self_node();
219 let outbox_store = connected.outbox_store;
220 let bootstrap_coordinator = connected.bootstrap_coordinator;
221 #[cfg(feature = "haematite-backend")]
222 let cluster_responder = connected.cluster_responder;
223 #[cfg(feature = "haematite-backend")]
224 let cluster_store = connected.cluster_store;
225 #[cfg(feature = "haematite-backend")]
226 let watched_peers = connected.watched_peers;
227 #[cfg(feature = "haematite-backend")]
231 let RoutingState {
232 shard_directory,
233 request_forwarder,
234 mint_routing,
235 } = build_routing_state(
236 cluster_store.as_ref(),
237 connected.directory_peers,
238 connected.self_node_id,
239 );
240 let (event_broadcast_capacity, query_timeout) = required_engine_seams(&runtime)?;
241 let (cluster_publisher, transcript_publisher) =
242 build_real_time_publishers(&runtime, connected.observability_store)?;
243 let metrics = Metrics::new().map_err(|error| metrics_config_error(&error))?;
244 let outbox_wake = Arc::new(tokio::sync::Notify::new());
249 let instrumented_store = Arc::new(
250 InstrumentedEventStore::new(
251 connected.event_store,
252 metrics.clone(),
253 runtime.default_namespace.clone(),
254 )
255 .with_outbox_wake(Arc::clone(&outbox_wake)),
256 );
257 let exported_metrics = runtime.metrics.enabled.then_some(metrics.clone());
258 #[cfg(not(feature = "haematite-backend"))]
259 let mint_routing: Option<crate::namespace::NamespaceRouting> = None;
260 let seams = build_worker_seams(
261 &runtime,
262 &cluster_publisher,
263 &connected.namespace_store,
264 &connected.worker_deployment_store,
265 mint_routing,
266 );
267 let (activity_dispatcher, activity_mock_registry, attempt_owners, workspace_root) =
268 build_decorated_dispatcher(&runtime, &seams, transcript_publisher.clone());
269
270 let engine = build_engine(EngineAssembly {
271 instrumented_store: &instrumented_store,
272 event_broadcast_capacity,
273 query_timeout,
274 activity_dispatcher,
275 active_registry: Arc::new(aion::Registry::default()),
276 bootstrap_coordinator,
277 runtime: &runtime,
278 })
279 .await?;
280 let engine = Arc::new(engine);
281 install_engine_backed_seams(&seams, &engine, runtime.outbox.enabled);
282 let resolver = NamespaceResolver::from_config(runtime.namespace.clone(), engine);
283 let worker_supervisor =
284 new_supervisor(&connected.worker_deployment_store, &cluster_publisher);
285 #[cfg(feature = "auth")]
286 let jwks_cache = build_jwks_cache(&runtime).await?;
287 Ok(Self {
288 inner: Arc::new(ServerStateInner {
289 namespace_guard: NamespaceGuard::new(resolver),
290 runtime,
291 metrics: exported_metrics,
292 worker_registry: seams.worker_registry,
293 pending_activities: seams.pending_activities,
294 heartbeat_tracker: seams.heartbeat_tracker,
295 drain_state: seams.drain_state,
296 health: Some(HealthState::new(instrumented_store, true)),
297 activity_mock_registry,
298 outbox_store,
299 namespace_store: connected.namespace_store,
300 worker_supervisor,
301 worker_deployment_store: connected.worker_deployment_store,
302 outbox_wake,
303 cluster_publisher,
304 transcript_publisher,
305 attempt_owners,
306 queue_service_state: seams.queue_service_state,
307 queue_declarations: seams.queue_declarations,
308 workspace_root,
309 cluster_self_node,
310 #[cfg(feature = "haematite-backend")]
311 cluster_responder,
312 #[cfg(feature = "haematite-backend")]
313 cluster_store,
314 #[cfg(feature = "haematite-backend")]
315 watched_peers,
316 #[cfg(feature = "haematite-backend")]
317 shard_directory,
318 #[cfg(feature = "haematite-backend")]
319 request_forwarder,
320 #[cfg(feature = "auth")]
321 jwks_cache,
322 }),
323 })
324 }
325
326 #[must_use]
328 pub fn from_parts(namespace_resolver: NamespaceResolver, runtime: RuntimeConfig) -> Self {
329 Self::from_parts_with_namespace_store(
334 namespace_resolver,
335 runtime,
336 Arc::new(aion_store::InMemoryStore::default()),
337 )
338 }
339
340 #[must_use]
347 pub fn from_parts_with_namespace_store<S>(
348 namespace_resolver: NamespaceResolver,
349 runtime: RuntimeConfig,
350 store: Arc<S>,
351 ) -> Self
352 where
353 S: NamespaceStore + WorkerDeploymentStore,
354 {
355 let namespace_store: Arc<dyn NamespaceStore> = store.clone();
356 let worker_deployment_store: Arc<dyn WorkerDeploymentStore> = store;
357 Self::from_parts_with_control_stores(
358 namespace_resolver,
359 runtime,
360 namespace_store,
361 worker_deployment_store,
362 )
363 }
364
365 #[must_use]
370 pub fn from_parts_with_control_stores(
371 namespace_resolver: NamespaceResolver,
372 runtime: RuntimeConfig,
373 namespace_store: Arc<dyn NamespaceStore>,
374 worker_deployment_store: Arc<dyn WorkerDeploymentStore>,
375 ) -> Self {
376 let heartbeat_tracker = HeartbeatTracker::new(runtime.worker.heartbeat_window);
377 let pending_activities =
384 PendingActivities::default().with_heartbeat_window(runtime.worker.heartbeat_window);
385 let bounds = transcript_bounds(&runtime);
389 let cluster_publisher = crate::cluster_publisher::ClusterEventPublisher::new(
390 Self::FALLBACK_CLUSTER_BROADCAST_CAPACITY,
391 );
392 Self {
393 inner: Arc::new(ServerStateInner {
394 namespace_guard: NamespaceGuard::new(namespace_resolver),
395 runtime,
396 worker_registry: ConnectedWorkerRegistry::default()
397 .with_worker_deployment_store(worker_deployment_store.clone())
398 .with_cluster_publisher(cluster_publisher.clone()),
399 pending_activities,
400 heartbeat_tracker,
401 drain_state: DrainState::default(),
402 metrics: None,
403 health: None,
404 activity_mock_registry: None,
405 outbox_store: None,
406 namespace_store,
407 worker_supervisor: new_supervisor(&worker_deployment_store, &cluster_publisher),
408 worker_deployment_store,
409 outbox_wake: Arc::new(tokio::sync::Notify::new()),
410 cluster_publisher,
411 transcript_publisher: build_transcript_publisher(
415 None,
416 Self::FALLBACK_CLUSTER_BROADCAST_CAPACITY,
417 bounds,
418 ),
419 attempt_owners: crate::worker::AttemptOwnerIndex::new(),
420 queue_service_state: crate::worker::QueueServiceState::default(),
421 queue_declarations: crate::worker::QueueDeclarationSource::default(),
422 workspace_root: crate::worker::WorkspaceRoot::resolve(),
426 cluster_self_node: None,
427 #[cfg(feature = "haematite-backend")]
428 cluster_responder: None,
429 #[cfg(feature = "haematite-backend")]
430 cluster_store: None,
431 #[cfg(feature = "haematite-backend")]
432 watched_peers: Vec::new(),
433 #[cfg(feature = "haematite-backend")]
434 shard_directory: None,
435 #[cfg(feature = "haematite-backend")]
436 request_forwarder: None,
437 #[cfg(feature = "auth")]
438 jwks_cache: None,
439 }),
440 }
441 }
442
443 #[cfg(feature = "auth")]
452 #[must_use]
453 pub fn from_parts_with_namespace_store_and_jwks<S>(
454 namespace_resolver: NamespaceResolver,
455 runtime: RuntimeConfig,
456 store: Arc<S>,
457 jwks_cache: JwksCache,
458 ) -> Self
459 where
460 S: NamespaceStore + WorkerDeploymentStore,
461 {
462 let namespace_store: Arc<dyn NamespaceStore> = store.clone();
463 let worker_deployment_store: Arc<dyn WorkerDeploymentStore> = store;
464 let heartbeat_tracker = HeartbeatTracker::new(runtime.worker.heartbeat_window);
465 let pending_activities =
472 PendingActivities::default().with_heartbeat_window(runtime.worker.heartbeat_window);
473 let bounds = transcript_bounds(&runtime);
477 let cluster_publisher = crate::cluster_publisher::ClusterEventPublisher::new(
478 Self::FALLBACK_CLUSTER_BROADCAST_CAPACITY,
479 );
480 Self {
481 inner: Arc::new(ServerStateInner {
482 namespace_guard: NamespaceGuard::new(namespace_resolver),
483 runtime,
484 worker_registry: ConnectedWorkerRegistry::default()
485 .with_worker_deployment_store(worker_deployment_store.clone())
486 .with_cluster_publisher(cluster_publisher.clone()),
487 pending_activities,
488 heartbeat_tracker,
489 drain_state: DrainState::default(),
490 metrics: None,
491 health: None,
492 activity_mock_registry: None,
493 outbox_store: None,
494 namespace_store,
495 worker_supervisor: new_supervisor(&worker_deployment_store, &cluster_publisher),
496 worker_deployment_store,
497 outbox_wake: Arc::new(tokio::sync::Notify::new()),
498 cluster_publisher,
499 transcript_publisher: build_transcript_publisher(
503 None,
504 Self::FALLBACK_CLUSTER_BROADCAST_CAPACITY,
505 bounds,
506 ),
507 attempt_owners: crate::worker::AttemptOwnerIndex::new(),
508 queue_service_state: crate::worker::QueueServiceState::default(),
509 queue_declarations: crate::worker::QueueDeclarationSource::default(),
510 workspace_root: crate::worker::WorkspaceRoot::resolve(),
514 cluster_self_node: None,
515 #[cfg(feature = "haematite-backend")]
516 cluster_responder: None,
517 #[cfg(feature = "haematite-backend")]
518 cluster_store: None,
519 #[cfg(feature = "haematite-backend")]
520 watched_peers: Vec::new(),
521 #[cfg(feature = "haematite-backend")]
522 shard_directory: None,
523 #[cfg(feature = "haematite-backend")]
524 request_forwarder: None,
525 jwks_cache: Some(jwks_cache),
526 }),
527 }
528 }
529
530 #[cfg(feature = "auth")]
536 #[must_use]
537 pub fn from_parts_with_jwks(
538 namespace_resolver: NamespaceResolver,
539 runtime: RuntimeConfig,
540 jwks_cache: JwksCache,
541 ) -> Self {
542 let fallback_deployment_store: Arc<dyn WorkerDeploymentStore> =
546 Arc::new(aion_store::InMemoryStore::default());
547 let heartbeat_tracker = HeartbeatTracker::new(runtime.worker.heartbeat_window);
548 let pending_activities =
555 PendingActivities::default().with_heartbeat_window(runtime.worker.heartbeat_window);
556 let bounds = transcript_bounds(&runtime);
560 let cluster_publisher = crate::cluster_publisher::ClusterEventPublisher::new(
561 Self::FALLBACK_CLUSTER_BROADCAST_CAPACITY,
562 );
563 Self {
564 inner: Arc::new(ServerStateInner {
565 namespace_guard: NamespaceGuard::new(namespace_resolver),
566 runtime,
567 worker_registry: ConnectedWorkerRegistry::default(),
568 pending_activities,
569 heartbeat_tracker,
570 drain_state: DrainState::default(),
571 metrics: None,
572 health: None,
573 activity_mock_registry: None,
574 outbox_store: None,
575 namespace_store: Arc::new(aion_store::InMemoryStore::default()),
580 worker_supervisor: new_supervisor(&fallback_deployment_store, &cluster_publisher),
581 worker_deployment_store: fallback_deployment_store,
582 outbox_wake: Arc::new(tokio::sync::Notify::new()),
583 cluster_publisher,
584 transcript_publisher: build_transcript_publisher(
588 None,
589 Self::FALLBACK_CLUSTER_BROADCAST_CAPACITY,
590 bounds,
591 ),
592 attempt_owners: crate::worker::AttemptOwnerIndex::new(),
593 queue_service_state: crate::worker::QueueServiceState::default(),
594 queue_declarations: crate::worker::QueueDeclarationSource::default(),
595 workspace_root: crate::worker::WorkspaceRoot::resolve(),
599 cluster_self_node: None,
600 #[cfg(feature = "haematite-backend")]
601 cluster_responder: None,
602 #[cfg(feature = "haematite-backend")]
603 cluster_store: None,
604 #[cfg(feature = "haematite-backend")]
605 watched_peers: Vec::new(),
606 #[cfg(feature = "haematite-backend")]
607 shard_directory: None,
608 #[cfg(feature = "haematite-backend")]
609 request_forwarder: None,
610 jwks_cache: Some(jwks_cache),
611 }),
612 }
613 }
614
615 #[must_use]
617 pub fn from_parts_with_registry(
618 namespace_resolver: NamespaceResolver,
619 runtime: RuntimeConfig,
620 worker_registry: ConnectedWorkerRegistry,
621 ) -> Self {
622 let fallback_deployment_store: Arc<dyn WorkerDeploymentStore> =
626 Arc::new(aion_store::InMemoryStore::default());
627 let heartbeat_tracker = HeartbeatTracker::new(runtime.worker.heartbeat_window);
628 let pending_activities =
635 PendingActivities::default().with_heartbeat_window(runtime.worker.heartbeat_window);
636 let bounds = transcript_bounds(&runtime);
640 let cluster_publisher = crate::cluster_publisher::ClusterEventPublisher::new(
641 Self::FALLBACK_CLUSTER_BROADCAST_CAPACITY,
642 );
643 Self {
644 inner: Arc::new(ServerStateInner {
645 namespace_guard: NamespaceGuard::new(namespace_resolver),
646 runtime,
647 worker_registry,
648 pending_activities,
649 heartbeat_tracker,
650 drain_state: DrainState::default(),
651 metrics: None,
652 health: None,
653 activity_mock_registry: None,
654 outbox_store: None,
655 namespace_store: Arc::new(aion_store::InMemoryStore::default()),
660 worker_supervisor: new_supervisor(&fallback_deployment_store, &cluster_publisher),
661 worker_deployment_store: fallback_deployment_store,
662 outbox_wake: Arc::new(tokio::sync::Notify::new()),
663 cluster_publisher,
664 transcript_publisher: build_transcript_publisher(
668 None,
669 Self::FALLBACK_CLUSTER_BROADCAST_CAPACITY,
670 bounds,
671 ),
672 attempt_owners: crate::worker::AttemptOwnerIndex::new(),
673 queue_service_state: crate::worker::QueueServiceState::default(),
674 queue_declarations: crate::worker::QueueDeclarationSource::default(),
675 workspace_root: crate::worker::WorkspaceRoot::resolve(),
679 cluster_self_node: None,
680 #[cfg(feature = "haematite-backend")]
681 cluster_responder: None,
682 #[cfg(feature = "haematite-backend")]
683 cluster_store: None,
684 #[cfg(feature = "haematite-backend")]
685 watched_peers: Vec::new(),
686 #[cfg(feature = "haematite-backend")]
687 shard_directory: None,
688 #[cfg(feature = "haematite-backend")]
689 request_forwarder: None,
690 #[cfg(feature = "auth")]
691 jwks_cache: None,
692 }),
693 }
694 }
695
696 #[must_use]
698 pub fn namespace_guard(&self) -> &NamespaceGuard {
699 &self.inner.namespace_guard
700 }
701
702 #[must_use]
704 pub fn deploy_guard(&self) -> crate::deploy::DeployGuard {
705 crate::deploy::DeployGuard::new(self.inner.namespace_guard.resolver().clone())
706 }
707
708 #[must_use]
710 pub fn runtime_config(&self) -> &RuntimeConfig {
711 &self.inner.runtime
712 }
713
714 #[must_use]
724 pub fn workspace_root(&self) -> &crate::worker::WorkspaceRoot {
725 &self.inner.workspace_root
726 }
727
728 #[must_use]
730 pub fn worker_registry(&self) -> &ConnectedWorkerRegistry {
731 &self.inner.worker_registry
732 }
733
734 #[must_use]
738 pub fn cluster_publisher(&self) -> &crate::cluster_publisher::ClusterEventPublisher {
739 &self.inner.cluster_publisher
740 }
741
742 #[must_use]
747 pub fn transcript_publisher(&self) -> &crate::activity_publisher::ActivityEventPublisher {
748 &self.inner.transcript_publisher
749 }
750
751 #[must_use]
755 pub fn attempt_owners(&self) -> &crate::worker::AttemptOwnerIndex {
756 &self.inner.attempt_owners
757 }
758
759 #[must_use]
762 pub fn queue_service_state(&self) -> &crate::worker::QueueServiceState {
763 &self.inner.queue_service_state
764 }
765
766 #[must_use]
770 pub fn queue_declarations(&self) -> &crate::worker::QueueDeclarationSource {
771 &self.inner.queue_declarations
772 }
773
774 pub fn unserved_queues(&self) -> Result<Vec<crate::worker::UnservedQueue>, ServerError> {
785 self.inner.queue_service_state.unserved()
786 }
787
788 pub fn unrecoverable_runs(
806 &self,
807 ) -> Result<Vec<(aion_core::WorkflowId, aion::registry::UnrecoverableRun)>, ServerError> {
808 self.engine()?
809 .registry()
810 .unrecoverable()
811 .list()
812 .map_err(ServerError::from)
813 }
814
815 #[must_use]
827 pub fn intervention_router(&self) -> crate::worker::InterventionRouter {
828 let transport: std::sync::Arc<dyn crate::worker::InterventionTransport> = {
829 #[cfg(feature = "liminal-transport")]
830 {
831 std::sync::Arc::new(crate::worker::LiminalInterventionTransport)
832 }
833 #[cfg(not(feature = "liminal-transport"))]
834 {
835 std::sync::Arc::new(NullInterventionTransport)
836 }
837 };
838 crate::worker::InterventionRouter::new(
839 self.inner.worker_registry.clone(),
840 self.inner.attempt_owners.clone(),
841 transport,
842 )
843 .with_transcript_publisher(self.inner.transcript_publisher.clone())
846 }
847
848 #[must_use]
852 pub fn cluster_self_node(&self) -> Option<&str> {
853 self.inner.cluster_self_node.as_deref()
854 }
855
856 pub fn engine(&self) -> Result<Arc<aion::Engine>, ServerError> {
869 self.inner
870 .namespace_guard
871 .resolver()
872 .engine()
873 .map(Arc::clone)
874 }
875
876 #[must_use]
878 pub fn pending_activities(&self) -> &PendingActivities {
879 &self.inner.pending_activities
880 }
881
882 #[must_use]
884 pub fn heartbeat_tracker(&self) -> &HeartbeatTracker {
885 &self.inner.heartbeat_tracker
886 }
887
888 #[must_use]
890 pub fn drain_state(&self) -> &DrainState {
891 &self.inner.drain_state
892 }
893
894 #[must_use]
896 pub fn metrics(&self) -> Option<&Metrics> {
897 self.inner.metrics.as_ref()
898 }
899
900 #[must_use]
902 pub fn health(&self) -> Option<&HealthState> {
903 self.inner.health.as_ref()
904 }
905
906 #[must_use]
911 pub fn activity_mock_registry(&self) -> Option<&ActivityMockRegistry> {
912 self.inner.activity_mock_registry.as_ref()
913 }
914
915 #[must_use]
921 pub fn outbox_store(&self) -> Option<Arc<dyn OutboxStore>> {
922 self.inner.outbox_store.clone()
923 }
924
925 #[must_use]
934 pub fn namespace_store(&self) -> &Arc<dyn NamespaceStore> {
935 &self.inner.namespace_store
936 }
937
938 #[must_use]
940 pub fn worker_deployment_store(&self) -> &Arc<dyn WorkerDeploymentStore> {
941 &self.inner.worker_deployment_store
942 }
943
944 #[must_use]
951 pub fn worker_supervisor(&self) -> &Arc<WorkerSupervisor> {
952 &self.inner.worker_supervisor
953 }
954
955 #[must_use]
964 pub fn namespace_minter(&self) -> NamespaceMinter {
965 let minter = NamespaceMinter::new(
966 Arc::clone(&self.inner.namespace_store),
967 self.inner.runtime.auto_create,
968 )
969 .with_cluster_publisher(self.inner.cluster_publisher.clone());
974 match self.namespace_routing() {
979 Some(routing) => minter.with_routing(routing),
980 None => minter,
981 }
982 }
983
984 #[cfg(feature = "haematite-backend")]
994 #[must_use]
995 pub fn namespace_routing(&self) -> Option<crate::namespace::NamespaceRouting> {
996 build_namespace_routing(
997 self.cluster_store(),
998 self.shard_directory(),
999 self.request_forwarder(),
1000 )
1001 }
1002
1003 #[cfg(not(feature = "haematite-backend"))]
1005 #[must_use]
1006 pub const fn namespace_routing(&self) -> Option<crate::namespace::NamespaceRouting> {
1007 None
1008 }
1009
1010 #[must_use]
1016 pub fn outbox_wake(&self) -> Arc<tokio::sync::Notify> {
1017 Arc::clone(&self.inner.outbox_wake)
1018 }
1019
1020 #[cfg(feature = "haematite-backend")]
1026 #[must_use]
1027 pub fn is_clustered(&self) -> bool {
1028 self.inner.cluster_responder.is_some()
1029 }
1030
1031 #[cfg(feature = "haematite-backend")]
1036 #[must_use]
1037 pub fn cluster_store(&self) -> Option<&Arc<aion_store_haematite::HaematiteStore>> {
1038 self.inner.cluster_store.as_ref()
1039 }
1040
1041 #[cfg(feature = "haematite-backend")]
1045 #[must_use]
1046 pub fn shard_directory(&self) -> Option<&Arc<crate::routing::StaticShardDirectory>> {
1047 self.inner.shard_directory.as_ref()
1048 }
1049
1050 #[cfg(feature = "haematite-backend")]
1053 #[must_use]
1054 pub fn request_forwarder(&self) -> Option<&Arc<dyn crate::routing::RequestForwarder>> {
1055 self.inner.request_forwarder.as_ref()
1056 }
1057
1058 #[must_use]
1075 pub fn spawn_heartbeat_sweeper(
1076 &self,
1077 shutdown: tokio::sync::watch::Receiver<bool>,
1078 ) -> tokio::task::JoinHandle<()> {
1079 let sweeper = crate::worker::HeartbeatSweeper::new(
1080 self.inner.heartbeat_tracker.clone(),
1081 self.inner.worker_registry.clone(),
1082 self.inner.pending_activities.clone(),
1083 self.inner.drain_state.clone(),
1084 self.inner.runtime.worker.heartbeat_window,
1085 )
1086 .with_queue_state(self.inner.queue_service_state.clone());
1090 tokio::spawn(sweeper.run(shutdown))
1091 }
1092
1093 #[cfg(feature = "liminal-transport")]
1104 #[must_use]
1105 pub fn spawn_liminal_liveness_probe(
1106 &self,
1107 notifier: std::sync::Arc<crate::worker::LiminalConnectionNotifier>,
1108 shutdown: tokio::sync::watch::Receiver<bool>,
1109 ) -> tokio::task::JoinHandle<()> {
1110 let probe = crate::worker::LivenessProbe::new(
1111 notifier,
1112 self.inner.heartbeat_tracker.clone(),
1113 self.inner.worker_registry.clone(),
1114 self.inner.runtime.worker.heartbeat_window,
1115 );
1116 tokio::spawn(probe.run(shutdown))
1117 }
1118
1119 #[cfg(feature = "haematite-backend")]
1135 pub fn spawn_cluster_supervisor(
1136 &self,
1137 config: crate::cluster::SupervisorConfig,
1138 shutdown: tokio::sync::watch::Receiver<bool>,
1139 ) -> Result<bool, ServerError> {
1140 let Some(cluster_store) = self.inner.cluster_store.clone() else {
1141 return Ok(false);
1142 };
1143 if self.inner.watched_peers.is_empty() {
1144 return Ok(false);
1145 }
1146 let engine = Arc::clone(self.inner.namespace_guard.resolver().engine()?);
1147 let publisher = Arc::new(self.inner.cluster_publisher.clone());
1151 let self_node = self.inner.cluster_self_node.clone().unwrap_or_default();
1152 let adopter = Arc::new(crate::cluster::OutboxSettlingAdopter::new(
1158 engine,
1159 self.inner.outbox_store.clone(),
1160 ));
1161 let supervisor = crate::cluster::ClusterSupervisor::new(
1162 cluster_store,
1163 adopter,
1164 self.inner.watched_peers.clone(),
1165 config,
1166 )
1167 .with_publisher(publisher, self_node);
1168 if !supervisor.watches_any() {
1169 return Ok(false);
1170 }
1171 tokio::spawn(supervisor.run(shutdown));
1172 Ok(true)
1173 }
1174
1175 #[cfg(feature = "auth")]
1177 #[must_use]
1178 pub fn jwks_cache(&self) -> Option<&JwksCache> {
1179 self.inner.jwks_cache.as_ref()
1180 }
1181
1182 pub fn shutdown(&self) -> Result<(), ServerError> {
1189 self.inner.namespace_guard.resolver().shutdown_engine()
1190 }
1191}
1192
1193#[cfg(feature = "auth")]
1194async fn build_jwks_cache(runtime: &RuntimeConfig) -> Result<Option<JwksCache>, ServerError> {
1195 if !runtime.auth.enabled {
1196 return Ok(None);
1197 }
1198 let Some(url) = runtime.auth.jwks_url.clone() else {
1199 return Err(ServerError::Config {
1200 message: "auth.jwks_url must not be empty when auth.enabled is true".to_owned(),
1201 });
1202 };
1203 let interval = std::time::Duration::from_secs(runtime.auth.jwks_refresh_seconds);
1204 let cache = JwksCache::new(url, interval)
1205 .await
1206 .map_err(|error| ServerError::Config {
1207 message: format!("auth jwks initial fetch failed: {error}"),
1208 })?;
1209 Ok(Some(cache))
1210}
1211
1212fn metrics_config_error(error: &MetricsError) -> ServerError {
1213 ServerError::Config {
1214 message: error.to_string(),
1215 }
1216}
1217
1218struct EngineAssembly<'a> {
1220 instrumented_store: &'a Arc<InstrumentedEventStore>,
1222 event_broadcast_capacity: std::num::NonZeroUsize,
1224 query_timeout: std::time::Duration,
1226 activity_dispatcher: Arc<dyn ActivityDispatcher>,
1228 active_registry: Arc<aion::Registry>,
1230 bootstrap_coordinator: bool,
1232 runtime: &'a RuntimeConfig,
1234}
1235
1236async fn build_engine(assembly: EngineAssembly<'_>) -> Result<aion::Engine, ServerError> {
1243 let mut search_attribute_schema = aion_core::SearchAttributeSchema::new();
1244 search_attribute_schema
1245 .register(
1246 crate::namespace::NAMESPACE_ATTRIBUTE,
1247 aion_core::SearchAttributeType::String,
1248 )
1249 .map_err(|error| ServerError::Config {
1250 message: format!("failed to register namespace search attribute: {error}"),
1251 })?;
1252 search_attribute_schema
1253 .register(
1254 crate::namespace::TASK_QUEUE_ATTRIBUTE,
1255 aion_core::SearchAttributeType::String,
1256 )
1257 .map_err(|error| ServerError::Config {
1258 message: format!("failed to register task_queue search attribute: {error}"),
1259 })?;
1260 let runtime = assembly.runtime;
1261 let builder = EngineBuilder::new()
1262 .store_arc(assembly.instrumented_store.clone())
1263 .event_streaming(assembly.event_broadcast_capacity)
1264 .in_memory_visibility()
1265 .search_attribute_schema(search_attribute_schema)
1266 .scheduler_threads(runtime.scheduler_threads)
1267 .outbox_enabled(runtime.outbox.enabled)
1268 .activity_dispatcher(assembly.activity_dispatcher)
1269 .active_registry(assembly.active_registry)
1270 .production_recovery_seam()
1271 .signal_router_factory(|runtime: Arc<RuntimeHandle>, handoff| {
1272 Arc::new(ConcreteSignalRouter::new(runtime, handoff)) as Arc<dyn SignalRouter>
1273 })
1274 .query_timeout(assembly.query_timeout)
1275 .bootstrap_schedule_coordinator(assembly.bootstrap_coordinator)
1281 .load_workflow_sources(runtime.workflow_packages.iter().map(PathBuf::as_path));
1282 let builder = if runtime.owned_shards.is_empty() {
1288 builder
1289 } else {
1290 builder.owned_shards(runtime.owned_shards.iter().copied())
1291 };
1292 builder.build().await.map_err(ServerError::from)
1293}
1294
1295fn required_engine_seams(
1300 runtime: &RuntimeConfig,
1301) -> Result<(std::num::NonZeroUsize, std::time::Duration), ServerError> {
1302 let event_broadcast_capacity = runtime
1303 .websocket
1304 .event_broadcast_capacity
1305 .and_then(std::num::NonZeroUsize::new)
1306 .ok_or_else(|| ServerError::Config {
1307 message: crate::config::EVENT_BROADCAST_CAPACITY_REQUIRED.to_owned(),
1308 })?;
1309 let query_timeout = runtime
1310 .query_timeout
1311 .filter(|timeout| !timeout.is_zero())
1312 .ok_or_else(|| ServerError::Config {
1313 message: crate::config::QUERY_TIMEOUT_REQUIRED.to_owned(),
1314 })?;
1315 Ok((event_broadcast_capacity, query_timeout))
1316}
1317
1318fn install_outbox_delivery(
1325 pending_activities: &PendingActivities,
1326 engine: &Arc<aion::Engine>,
1327 outbox_enabled: bool,
1328) {
1329 if outbox_enabled {
1330 let callback = Arc::new(crate::worker::ServerOutboxDeliveryCallback::new(
1331 Arc::clone(engine),
1332 ));
1333 pending_activities.set_outbox_delivery(callback);
1334 }
1335}
1336
1337struct WorkerSeams {
1339 worker_registry: ConnectedWorkerRegistry,
1340 pending_activities: PendingActivities,
1341 heartbeat_tracker: HeartbeatTracker,
1342 drain_state: DrainState,
1343 queue_declarations: crate::worker::QueueDeclarationSource,
1344 queue_service_state: crate::worker::QueueServiceState,
1345 declared_bodies: crate::worker::DeclaredBodySource,
1346}
1347
1348fn build_worker_seams(
1353 runtime: &RuntimeConfig,
1354 cluster_publisher: &crate::cluster_publisher::ClusterEventPublisher,
1355 namespace_store: &Arc<dyn NamespaceStore>,
1356 worker_deployment_store: &Arc<dyn WorkerDeploymentStore>,
1357 mint_routing: Option<crate::namespace::NamespaceRouting>,
1358) -> WorkerSeams {
1359 let worker_registry = ConnectedWorkerRegistry::default()
1360 .with_cluster_publisher(cluster_publisher.clone())
1361 .with_namespace_minting(namespace_store.clone(), runtime.auto_create)
1362 .with_worker_deployment_store(worker_deployment_store.clone());
1363 let worker_registry = match mint_routing {
1369 Some(routing) => worker_registry.with_namespace_routing(routing),
1370 None => worker_registry,
1371 };
1372 WorkerSeams {
1373 worker_registry,
1374 pending_activities: PendingActivities::default()
1379 .with_heartbeat_window(runtime.worker.heartbeat_window),
1380 heartbeat_tracker: HeartbeatTracker::new(runtime.worker.heartbeat_window),
1381 drain_state: DrainState::default(),
1382 queue_declarations: crate::worker::QueueDeclarationSource::default(),
1383 queue_service_state: crate::worker::QueueServiceState::default(),
1384 declared_bodies: crate::worker::DeclaredBodySource::default(),
1385 }
1386}
1387
1388fn install_queue_declarations(
1395 queue_declarations: &crate::worker::QueueDeclarationSource,
1396 engine: &Arc<aion::Engine>,
1397) {
1398 queue_declarations.install(Arc::new(crate::worker::EngineQueueDeclarations::new(
1399 Arc::clone(engine),
1400 )));
1401}
1402
1403fn install_declared_bodies(
1411 declared_bodies: &crate::worker::DeclaredBodySource,
1412 engine: &Arc<aion::Engine>,
1413) {
1414 declared_bodies.install(Arc::new(crate::worker::EngineDeclaredBodies::new(
1415 Arc::clone(engine),
1416 )));
1417}
1418
1419fn install_engine_backed_seams(seams: &WorkerSeams, engine: &Arc<aion::Engine>, outbox: bool) {
1425 install_outbox_delivery(&seams.pending_activities, engine, outbox);
1426 install_queue_declarations(&seams.queue_declarations, engine);
1427 install_declared_bodies(&seams.declared_bodies, engine);
1428}
1429
1430fn build_bridge_dispatcher(
1443 runtime: &RuntimeConfig,
1444 seams: &WorkerSeams,
1445) -> (WorkerActivityDispatcher, crate::worker::AttemptOwnerIndex) {
1446 let attempt_owners = crate::worker::AttemptOwnerIndex::new();
1447 let dispatcher = WorkerActivityDispatcher::new(
1448 seams.worker_registry.clone(),
1449 runtime.default_namespace.clone(),
1450 seams.heartbeat_tracker.clone(),
1451 )
1452 .with_pending(seams.pending_activities.clone())
1453 .with_drain_state(seams.drain_state.clone())
1454 .with_tokio_handle(tokio::runtime::Handle::current())
1455 .with_attempt_owners(attempt_owners.clone())
1456 .with_queue_service(runtime.worker.queue_service.clone())
1457 .with_queue_declarations(seams.queue_declarations.clone())
1458 .with_queue_state(seams.queue_service_state.clone());
1459 (dispatcher, attempt_owners)
1460}
1461
1462fn build_decorated_dispatcher(
1465 runtime: &RuntimeConfig,
1466 seams: &WorkerSeams,
1467 transcript: crate::activity_publisher::ActivityEventPublisher,
1468) -> (
1469 Arc<dyn ActivityDispatcher>,
1470 Option<ActivityMockRegistry>,
1471 crate::worker::AttemptOwnerIndex,
1472 crate::worker::WorkspaceRoot,
1473) {
1474 let workspace_root = crate::worker::WorkspaceRoot::resolve();
1478 let (dispatcher, attempt_owners) = build_bridge_dispatcher(runtime, seams);
1479 let (activity_dispatcher, activity_mock_registry) = decorate_activity_dispatcher(
1480 dispatcher,
1481 seams.declared_bodies.clone(),
1482 workspace_root.clone(),
1483 transcript,
1484 runtime.dev.enabled,
1485 );
1486 (
1487 activity_dispatcher,
1488 activity_mock_registry,
1489 attempt_owners,
1490 workspace_root,
1491 )
1492}
1493
1494fn decorate_activity_dispatcher(
1495 dispatcher: WorkerActivityDispatcher,
1496 declared_bodies: crate::worker::DeclaredBodySource,
1497 workspace_root: crate::worker::WorkspaceRoot,
1498 transcript: crate::activity_publisher::ActivityEventPublisher,
1499 dev_enabled: bool,
1500) -> (Arc<dyn ActivityDispatcher>, Option<ActivityMockRegistry>) {
1501 let declared = crate::worker::DeclaredCommandDispatcher::new(
1507 Arc::new(dispatcher),
1508 declared_bodies,
1509 tokio::runtime::Handle::current(),
1510 workspace_root,
1511 transcript,
1512 );
1513 if dev_enabled {
1514 let registry = ActivityMockRegistry::new();
1515 let decorated = DevMockingDispatcher::new(Arc::new(declared), registry.clone());
1516 (Arc::new(decorated), Some(registry))
1517 } else {
1518 (Arc::new(declared), None)
1519 }
1520}
1521
1522fn required_cluster_broadcast_capacity(
1527 runtime: &RuntimeConfig,
1528) -> Result<std::num::NonZeroUsize, ServerError> {
1529 runtime
1530 .websocket
1531 .cluster_broadcast_capacity
1532 .and_then(std::num::NonZeroUsize::new)
1533 .ok_or_else(|| ServerError::Config {
1534 message: crate::config::CLUSTER_BROADCAST_CAPACITY_REQUIRED.to_owned(),
1535 })
1536}
1537
1538fn build_real_time_publishers(
1551 runtime: &RuntimeConfig,
1552 observability_store: Option<Arc<dyn aion_store::ObservabilityStore>>,
1553) -> Result<
1554 (
1555 crate::cluster_publisher::ClusterEventPublisher,
1556 crate::activity_publisher::ActivityEventPublisher,
1557 ),
1558 ServerError,
1559> {
1560 let capacity = required_cluster_broadcast_capacity(runtime)?;
1561 Ok((
1562 crate::cluster_publisher::ClusterEventPublisher::new(capacity),
1563 build_transcript_publisher(observability_store, capacity, transcript_bounds(runtime)),
1564 ))
1565}
1566
1567fn transcript_bounds(runtime: &RuntimeConfig) -> crate::activity_bounds::TranscriptBounds {
1569 crate::activity_bounds::TranscriptBounds {
1570 max_event_bytes: runtime.observability.max_event_bytes,
1571 max_stream_events: runtime.observability.max_stream_events,
1572 }
1573}
1574
1575fn build_transcript_publisher(
1586 observability_store: Option<Arc<dyn aion_store::ObservabilityStore>>,
1587 capacity: std::num::NonZeroUsize,
1588 bounds: crate::activity_bounds::TranscriptBounds,
1589) -> crate::activity_publisher::ActivityEventPublisher {
1590 let store = observability_store
1591 .unwrap_or_else(|| Arc::new(aion_store::InMemoryObservabilityStore::default()));
1592 crate::activity_publisher::ActivityEventPublisher::new(store, capacity).with_bounds(bounds)
1593}
1594
1595#[cfg(feature = "haematite-backend")]
1597struct RoutingState {
1598 shard_directory: Option<Arc<crate::routing::StaticShardDirectory>>,
1599 request_forwarder: Option<Arc<dyn crate::routing::RequestForwarder>>,
1600 mint_routing: Option<crate::namespace::NamespaceRouting>,
1605}
1606
1607#[cfg(feature = "haematite-backend")]
1611fn build_routing_state(
1612 cluster_store: Option<&Arc<aion_store_haematite::HaematiteStore>>,
1613 directory_peers: Vec<crate::routing::DirectoryPeer>,
1614 self_node_id: Option<String>,
1615) -> RoutingState {
1616 let Some(store) = cluster_store else {
1617 return RoutingState {
1618 shard_directory: None,
1619 request_forwarder: None,
1620 mint_routing: None,
1621 };
1622 };
1623 let shard_directory = Arc::new(crate::routing::StaticShardDirectory::new(
1624 Arc::clone(store),
1625 directory_peers,
1626 self_node_id,
1627 ));
1628 let request_forwarder: Arc<dyn crate::routing::RequestForwarder> =
1629 Arc::new(crate::routing::GrpcRequestForwarder::new());
1630 let mint_routing = build_namespace_routing(
1631 Some(store),
1632 Some(&shard_directory),
1633 Some(&request_forwarder),
1634 );
1635 RoutingState {
1636 shard_directory: Some(shard_directory),
1637 request_forwarder: Some(request_forwarder),
1638 mint_routing,
1639 }
1640}
1641
1642#[cfg(feature = "haematite-backend")]
1652fn build_namespace_routing(
1653 cluster_store: Option<&Arc<aion_store_haematite::HaematiteStore>>,
1654 shard_directory: Option<&Arc<crate::routing::StaticShardDirectory>>,
1655 request_forwarder: Option<&Arc<dyn crate::routing::RequestForwarder>>,
1656) -> Option<crate::namespace::NamespaceRouting> {
1657 use crate::namespace::{
1658 GrpcMintForwarder, MintForwarder, MintShardOwners, NamespaceRouting, NamespaceShardResolver,
1659 };
1660 let store = Arc::clone(cluster_store?);
1661 let directory = Arc::clone(shard_directory?);
1662 let shards: Arc<dyn NamespaceShardResolver> = store;
1663 let owners: Arc<dyn MintShardOwners> = directory;
1664 let forwarder: Arc<dyn MintForwarder> =
1665 Arc::new(GrpcMintForwarder::new(Arc::clone(request_forwarder?)));
1666 Some(NamespaceRouting::new(shards, owners, forwarder))
1667}
1668
1669struct ConnectedStore {
1679 event_store: Arc<dyn EventStore>,
1680 outbox_store: Option<Arc<dyn OutboxStore>>,
1681 namespace_store: Arc<dyn NamespaceStore>,
1688 worker_deployment_store: Arc<dyn WorkerDeploymentStore>,
1690 observability_store: Option<Arc<dyn aion_store::ObservabilityStore>>,
1698 bootstrap_coordinator: bool,
1699 #[cfg(feature = "haematite-backend")]
1700 cluster_responder: Option<aion_store_haematite::ClusterResponder>,
1701 #[cfg(feature = "haematite-backend")]
1705 cluster_store: Option<Arc<aion_store_haematite::HaematiteStore>>,
1706 #[cfg(feature = "haematite-backend")]
1709 watched_peers: Vec<crate::cluster::WatchedPeer>,
1710 #[cfg(feature = "haematite-backend")]
1714 directory_peers: Vec<crate::routing::DirectoryPeer>,
1715 #[cfg(feature = "haematite-backend")]
1719 self_node_id: Option<String>,
1720}
1721
1722impl ConnectedStore {
1723 fn local(
1730 event_store: Arc<dyn EventStore>,
1731 outbox_store: Option<Arc<dyn OutboxStore>>,
1732 namespace_store: Arc<dyn NamespaceStore>,
1733 worker_deployment_store: Arc<dyn WorkerDeploymentStore>,
1734 ) -> Self {
1735 Self {
1736 event_store,
1737 outbox_store,
1738 namespace_store,
1739 worker_deployment_store,
1740 observability_store: None,
1744 bootstrap_coordinator: true,
1745 #[cfg(feature = "haematite-backend")]
1746 cluster_responder: None,
1747 #[cfg(feature = "haematite-backend")]
1748 cluster_store: None,
1749 #[cfg(feature = "haematite-backend")]
1750 watched_peers: Vec::new(),
1751 #[cfg(feature = "haematite-backend")]
1752 directory_peers: Vec::new(),
1753 #[cfg(feature = "haematite-backend")]
1754 self_node_id: None,
1755 }
1756 }
1757
1758 #[cfg(feature = "haematite-backend")]
1761 fn cluster_self_node(&self) -> Option<String> {
1762 self.self_node_id.clone()
1763 }
1764
1765 #[cfg(not(feature = "haematite-backend"))]
1767 const fn cluster_self_node(&self) -> Option<String> {
1768 None
1769 }
1770}
1771
1772async fn connect_store(config: StoreConfig) -> Result<ConnectedStore, ServerError> {
1782 match config.backend {
1783 StoreBackend::Memory => {
1784 let leaf = Arc::new(aion_store::InMemoryStore::default());
1787 let namespace_store: Arc<dyn NamespaceStore> = leaf.clone();
1788 let worker_deployment_store: Arc<dyn WorkerDeploymentStore> = leaf.clone();
1789 Ok(ConnectedStore::local(
1790 leaf,
1791 None,
1792 namespace_store,
1793 worker_deployment_store,
1794 ))
1795 }
1796 StoreBackend::LibSql => {
1797 #[cfg(feature = "libsql-backend")]
1798 {
1799 connect_libsql_store(config).await
1800 }
1801 #[cfg(not(feature = "libsql-backend"))]
1802 {
1803 let _ = config;
1804 connect_libsql_store_unavailable()
1805 }
1806 }
1807 StoreBackend::Haematite => {
1808 #[cfg(feature = "haematite-backend")]
1809 {
1810 connect_haematite_store(config).await
1811 }
1812 #[cfg(not(feature = "haematite-backend"))]
1813 {
1814 let _ = config;
1815 connect_haematite_store_unavailable()
1816 }
1817 }
1818 }
1819}
1820
1821#[cfg(feature = "libsql-backend")]
1825async fn connect_libsql_store(config: StoreConfig) -> Result<ConnectedStore, ServerError> {
1826 let Some(url) = config.url else {
1827 return Err(ServerError::Config {
1828 message: "store.url must not be empty when store.backend is libsql".to_owned(),
1829 });
1830 };
1831 let store = LibSqlStore::open(url.clone())
1832 .await
1833 .map_err(ServerError::from)?;
1834 store
1835 .validate_event_compatibility()
1836 .await
1837 .map_err(|error| match error {
1838 aion_store::StoreError::Serialization(_) => ServerError::Config {
1839 message: format!(
1840 "Database schema mismatch — delete {url} and restart, or run migrations."
1841 ),
1842 },
1843 other => ServerError::from(other),
1844 })?;
1845 let leaf = Arc::new(store);
1846 let event_store: Arc<dyn EventStore> = leaf.clone();
1847 let namespace_store: Arc<dyn NamespaceStore> = leaf.clone();
1848 let worker_deployment_store: Arc<dyn WorkerDeploymentStore> = leaf.clone();
1849 let outbox_store: Arc<dyn OutboxStore> = leaf;
1850 Ok(ConnectedStore::local(
1851 event_store,
1852 Some(outbox_store),
1853 namespace_store,
1854 worker_deployment_store,
1855 ))
1856}
1857
1858#[cfg(not(feature = "libsql-backend"))]
1862fn connect_libsql_store_unavailable() -> Result<ConnectedStore, ServerError> {
1863 Err(ServerError::Config {
1864 message: "store.backend = libsql requires the aion-server `libsql-backend` feature"
1865 .to_owned(),
1866 })
1867}
1868
1869#[cfg(feature = "haematite-backend")]
1889async fn connect_haematite_store(config: StoreConfig) -> Result<ConnectedStore, ServerError> {
1890 let Some(data_dir) = config.data_dir else {
1891 return Err(ServerError::Config {
1892 message: "store.data_dir must not be empty when store.backend is haematite".to_owned(),
1893 });
1894 };
1895 let shard_count = config.shard_count;
1896 let owned_shards = config.owned_shards.clone();
1897 let cluster = config.cluster.clone();
1898 let watched_peers: Vec<crate::cluster::WatchedPeer> = cluster
1903 .as_ref()
1904 .map(|cluster| {
1905 cluster
1906 .peers
1907 .iter()
1908 .map(|peer| crate::cluster::WatchedPeer {
1909 name: peer.name.clone(),
1910 owned_shards: peer.owned_shards.clone(),
1911 })
1912 .collect()
1913 })
1914 .unwrap_or_default();
1915 let directory_peers: Vec<crate::routing::DirectoryPeer> = cluster
1918 .as_ref()
1919 .map(|cluster| {
1920 cluster
1921 .peers
1922 .iter()
1923 .map(|peer| crate::routing::DirectoryPeer {
1924 name: peer.name.clone(),
1925 owned_shards: peer.owned_shards.clone(),
1926 grpc_addr: peer.grpc_address,
1927 })
1928 .collect()
1929 })
1930 .unwrap_or_default();
1931 let self_node_id: Option<String> = cluster.as_ref().map(|cluster| cluster.node_id.clone());
1934 let (store, responder) =
1938 tokio::task::spawn_blocking(move || build_haematite_store(&data_dir, shard_count, cluster))
1939 .await
1940 .map_err(|error| ServerError::Config {
1941 message: format!("haematite store initialization task failed: {error}"),
1942 })??;
1943
1944 let bootstrap_coordinator = if owned_shards.is_empty() {
1948 true
1949 } else {
1950 store.set_owned_shards(owned_shards.iter().copied());
1951 store.owns_workflow_shard(&aion::schedule_coordinator_workflow_id())
1952 };
1953
1954 let leaf = Arc::new(store);
1955 let event_store: Arc<dyn EventStore> = leaf.clone();
1956 let outbox_store: Arc<dyn OutboxStore> = leaf.clone();
1957 let namespace_store: Arc<dyn NamespaceStore> = leaf.clone();
1961 let worker_deployment_store: Arc<dyn WorkerDeploymentStore> = leaf.clone();
1962 let observability_store: Arc<dyn aion_store::ObservabilityStore> = leaf.clone();
1967 let cluster_store = responder.as_ref().map(|_| leaf);
1971 let (watched_peers, directory_peers, self_node_id) = if cluster_store.is_some() {
1972 (watched_peers, directory_peers, self_node_id)
1973 } else {
1974 (Vec::new(), Vec::new(), None)
1975 };
1976 Ok(ConnectedStore {
1977 event_store,
1978 outbox_store: Some(outbox_store),
1979 namespace_store,
1980 worker_deployment_store,
1981 observability_store: Some(observability_store),
1982 bootstrap_coordinator,
1983 cluster_responder: responder,
1984 cluster_store,
1985 watched_peers,
1986 directory_peers,
1987 self_node_id,
1988 })
1989}
1990
1991#[cfg(feature = "haematite-backend")]
2006fn build_haematite_store(
2007 data_dir: &str,
2008 shard_count: usize,
2009 cluster: Option<crate::config::ClusterConfig>,
2010) -> Result<
2011 (
2012 aion_store_haematite::HaematiteStore,
2013 Option<aion_store_haematite::ClusterResponder>,
2014 ),
2015 ServerError,
2016> {
2017 build_haematite_store_with_hook(data_dir, shard_count, cluster, || Ok(()))
2018}
2019
2020#[cfg(feature = "haematite-backend")]
2021fn build_haematite_store_with_hook(
2022 data_dir: &str,
2023 shard_count: usize,
2024 cluster: Option<crate::config::ClusterConfig>,
2025 before_backend_touch: impl FnOnce() -> Result<(), std::io::Error>,
2026) -> Result<
2027 (
2028 aion_store_haematite::HaematiteStore,
2029 Option<aion_store_haematite::ClusterResponder>,
2030 ),
2031 ServerError,
2032> {
2033 use aion_store_haematite::{ClusterBootstrap, HaematiteStore};
2034
2035 let private_root = crate::filesystem::ConfinedDir::open_or_create(std::path::Path::new(
2043 data_dir,
2044 ))
2045 .map_err(|error| ServerError::Config {
2046 message: format!("unsafe store.data_dir `{data_dir}`: {error}"),
2047 })?;
2048
2049 for shard in 0..shard_count {
2053 private_root
2054 .create_dir_all(std::path::Path::new(&format!("shard-{shard}")))
2055 .map_err(|error| ServerError::Config {
2056 message: format!(
2057 "failed to materialize shard-{shard} under store.data_dir `{data_dir}`: {error}"
2058 ),
2059 })?;
2060 }
2061 private_root
2062 .harden_tree()
2063 .map_err(|error| private_store_mode_error(data_dir, &error))?;
2064
2065 before_backend_touch().map_err(|error| ServerError::Config {
2068 message: format!("store.data_dir pre-open hook failed: {error}"),
2069 })?;
2070
2071 #[cfg(unix)]
2072 let backend_path = private_root
2073 .backend_path()
2074 .map_err(|error| ServerError::Config {
2075 message: format!("failed to resolve held store.data_dir `{data_dir}`: {error}"),
2076 })?;
2077 #[cfg(all(unix, not(any(target_os = "linux", target_os = "android"))))]
2078 crate::filesystem::validate_ambient_backend_ancestors(&backend_path).map_err(|error| {
2079 let (component, reason) = error.into_parts();
2080 ServerError::UnsafeDataRootAncestor {
2081 data_root: backend_path.clone(),
2082 component,
2083 reason,
2084 }
2085 })?;
2086 #[cfg(not(unix))]
2087 let backend_path = std::path::PathBuf::from(data_dir);
2088
2089 let Some(cluster) = cluster else {
2090 let store = if backend_path.join("config.json").exists() {
2091 HaematiteStore::open(&backend_path).map_err(ServerError::from)?
2092 } else {
2093 HaematiteStore::create_with_shard_count(&backend_path, shard_count)
2094 .map_err(ServerError::from)?
2095 };
2096 store.materialize_all_shards().map_err(ServerError::from)?;
2097 private_root
2098 .harden_tree()
2099 .map_err(|error| private_store_mode_error(data_dir, &error))?;
2100 let store = store.retain_data_root_capability(private_root);
2101 return Ok((store, None));
2102 };
2103
2104 let boot = ClusterBootstrap {
2105 node_id: cluster.node_id,
2106 bind_address: cluster.bind_address,
2107 members: cluster.members,
2108 peers: cluster
2109 .peers
2110 .into_iter()
2111 .map(|peer| (peer.name, peer.address))
2112 .collect(),
2113 timeout: HAEMATITE_CLUSTER_OP_TIMEOUT,
2114 };
2115 let (store, responder) =
2116 HaematiteStore::open_or_create_distributed(&backend_path, shard_count, boot)
2117 .map_err(ServerError::from)?;
2118 store.materialize_all_shards().map_err(ServerError::from)?;
2119 private_root
2120 .harden_tree()
2121 .map_err(|error| private_store_mode_error(data_dir, &error))?;
2122 let store = store.retain_data_root_capability(private_root);
2123 Ok((store, Some(responder)))
2124}
2125
2126#[cfg(feature = "haematite-backend")]
2127fn private_store_mode_error(data_dir: &str, error: &std::io::Error) -> ServerError {
2128 ServerError::Config {
2129 message: format!(
2130 "failed to apply private modes under store.data_dir `{data_dir}`: {error}"
2131 ),
2132 }
2133}
2134
2135#[cfg(feature = "haematite-backend")]
2137const HAEMATITE_CLUSTER_OP_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(5);
2138
2139#[cfg(not(feature = "haematite-backend"))]
2143fn connect_haematite_store_unavailable() -> Result<ConnectedStore, ServerError> {
2144 Err(ServerError::Config {
2145 message: "store.backend = haematite requires the aion-server `haematite-backend` feature"
2146 .to_owned(),
2147 })
2148}
2149
2150#[cfg(not(feature = "liminal-transport"))]
2159#[derive(Clone, Debug)]
2160struct NullInterventionTransport;
2161
2162#[cfg(not(feature = "liminal-transport"))]
2163#[async_trait::async_trait]
2164impl crate::worker::InterventionTransport for NullInterventionTransport {
2165 async fn push(
2166 &self,
2167 _worker: &crate::worker::WorkerHandle,
2168 _command: aion_core::InterventionCommand,
2169 ) -> Result<aion_core::InterventionOutcome, ServerError> {
2170 Err(ServerError::worker_connection_lost(
2171 "intervention",
2172 "no intervention push transport is compiled in".to_owned(),
2173 ))
2174 }
2175}
2176
2177#[cfg(test)]
2178mod tests {
2179 use std::{net::SocketAddr, time::Duration};
2180
2181 use aion_store::InMemoryStore;
2182
2183 use super::ServerState;
2184 use crate::config::{
2185 AuthConfig, AuthoringConfig, DeployConfig, DevConfig, ListenConfig, MetricsConfig,
2186 NamespaceConfig, NamespaceMode, OpsConsoleAssetSource, OpsConsoleConfig, OutboxConfig,
2187 RuntimeConfig, WebSocketConfig, WorkerConfig,
2188 };
2189
2190 fn runtime_config() -> RuntimeConfig {
2191 RuntimeConfig {
2192 listen: ListenConfig {
2193 grpc: SocketAddr::from(([127, 0, 0, 1], 50051)),
2194 http: SocketAddr::from(([127, 0, 0, 1], 8080)),
2195 },
2196 tls: None,
2197 auth: AuthConfig {
2198 enabled: false,
2199 jwks_url: None,
2200 jwks_refresh_seconds: 300,
2201 },
2202 ops_console: OpsConsoleConfig {
2203 source: OpsConsoleAssetSource::Embedded,
2204 },
2205 namespace: NamespaceConfig {
2206 mode: NamespaceMode::SharedEngine,
2207 },
2208 worker: WorkerConfig {
2209 heartbeat_window: Duration::from_secs(30),
2210 ..WorkerConfig::default()
2211 },
2212 websocket: WebSocketConfig {
2213 outbound_buffer_bound: 32,
2214 event_broadcast_capacity: Some(64),
2215 cluster_broadcast_capacity: Some(64),
2216 },
2217 workflow_packages: Vec::new(),
2218 deploy: DeployConfig::default(),
2219 authoring: AuthoringConfig::default(),
2220 dev: DevConfig::default(),
2221 outbox: OutboxConfig::default(),
2222 observability: crate::config::ObservabilityConfig::default(),
2223 mcp: crate::config::ResolvedMcpConfig::default(),
2224 scheduler_threads: 1,
2225 query_timeout: Some(Duration::from_secs(10)),
2226 default_namespace: "default".to_owned(),
2227 auto_create: crate::config::AutoCreate::Open,
2228 max_in_flight_activities: crate::config::DEFAULT_MAX_IN_FLIGHT_ACTIVITIES,
2229 drain_timeout: Duration::from_secs(30),
2230 metrics: MetricsConfig { enabled: true },
2231 owned_shards: Vec::new(),
2232 cors_allowed_origins: Vec::new(),
2233 }
2234 }
2235
2236 #[tokio::test]
2237 async fn builds_state_with_in_memory_store() -> Result<(), Box<dyn std::error::Error>> {
2238 let state =
2239 ServerState::build_with_store(InMemoryStore::default(), runtime_config()).await?;
2240
2241 std::hint::black_box(state.namespace_guard());
2242 std::hint::black_box(state.worker_registry());
2243
2244 Ok(())
2245 }
2246
2247 #[tokio::test]
2255 async fn unserved_queues_surfaces_a_parked_dispatch_and_clears_it()
2256 -> Result<(), Box<dyn std::error::Error>> {
2257 use aion::{ActivityDispatch, ActivityDispatcher as _};
2258 use aion_core::{ActivityId, RunId, WorkflowId};
2259 use std::sync::Arc;
2260
2261 let state =
2262 ServerState::build_with_store(InMemoryStore::default(), runtime_config()).await?;
2263 assert!(
2264 state.unserved_queues()?.is_empty(),
2265 "a calm boot has no unserved queues"
2266 );
2267 assert!(state.queue_declarations().is_installed());
2271 assert_eq!(
2272 state
2273 .queue_declarations()
2274 .declaration_for("nobody-serves-this"),
2275 crate::worker::QueueDeclaration::Unknown
2276 );
2277
2278 let dispatcher = Arc::new(
2279 crate::worker::WorkerActivityDispatcher::new(
2280 state.worker_registry().clone(),
2281 "default",
2282 crate::worker::HeartbeatTracker::new(Duration::from_secs(5)),
2283 )
2284 .with_queue_state(state.queue_service_state().clone())
2285 .with_queue_declarations(state.queue_declarations().clone()),
2286 );
2287 let workflow_id = WorkflowId::new_v4();
2288 let request = ActivityDispatch {
2289 namespace: "default".to_owned(),
2290 task_queue: "nobody-serves-this".to_owned(),
2291 node: None,
2292 workflow_id: workflow_id.clone(),
2293 run_id: RunId::new_v4(),
2294 activity_id: ActivityId::from_sequence_position(0),
2295 name: "greet".to_owned(),
2296 input: "{}".to_owned(),
2297 config: "{}".to_owned(),
2298 attempt: 1,
2299 advisory: false,
2300 labels: std::collections::BTreeMap::new(),
2301 };
2302 let parked = std::thread::spawn(move || dispatcher.dispatch(request));
2303
2304 let mut unserved = Vec::new();
2305 for _ in 0..30 {
2306 unserved = state.unserved_queues()?;
2307 if !unserved.is_empty() {
2308 break;
2309 }
2310 tokio::time::sleep(Duration::from_millis(100)).await;
2311 }
2312 assert_eq!(unserved.len(), 1, "the parked dispatch is not surfaced");
2313 assert_eq!(
2314 unserved[0].reason,
2315 crate::worker::QueueServiceReason::NoLivePollers,
2316 "an empty catalog must not be read as a structural refusal"
2317 );
2318 assert_eq!(unserved[0].key.task_queue, "nobody-serves-this");
2319 assert_eq!(unserved[0].waiting.len(), 1);
2320 assert_eq!(unserved[0].waiting[0].workflow_id, workflow_id);
2321
2322 let (worker_tx, worker_rx) = tokio::sync::mpsc::channel(1);
2324 drop(worker_rx);
2325 let registration = state.worker_registry().register_namespaces(
2326 [String::from("default")],
2327 "nobody-serves-this",
2328 None,
2329 [String::from("greet")].iter(),
2330 worker_tx,
2331 )?;
2332 let outcome = parked.join().map_err(|_| "parked dispatch panicked")?;
2333 assert!(outcome.is_err(), "the released dispatch must resolve");
2334 assert!(
2335 state.unserved_queues()?.is_empty(),
2336 "a resolved dispatch must leave the unserved state"
2337 );
2338 registration.deregister()?;
2339 Ok(())
2340 }
2341
2342 #[tokio::test]
2343 async fn namespace_store_is_reachable_and_functional_after_default_boot()
2344 -> Result<(), Box<dyn std::error::Error>> {
2345 use aion_store::{MintOutcome, NamespaceOrigin};
2346
2347 let state =
2351 ServerState::build_with_store(InMemoryStore::default(), runtime_config()).await?;
2352
2353 let store = state.namespace_store();
2354
2355 let outcome = store
2357 .register_namespace("orders", NamespaceOrigin::WorkerMint)
2358 .await?;
2359 assert_eq!(
2360 outcome,
2361 MintOutcome::Created,
2362 "the first reference to a namespace mints it"
2363 );
2364
2365 let again = store
2367 .register_namespace("orders", NamespaceOrigin::WorkerMint)
2368 .await?;
2369 assert_eq!(
2370 again,
2371 MintOutcome::AlreadyExisted,
2372 "a second reference touches the existing record rather than re-creating it"
2373 );
2374
2375 let fetched = store.get_namespace("orders").await?;
2377 let record = fetched.ok_or("registered namespace must be retrievable via get_namespace")?;
2378 assert_eq!(record.name, "orders");
2379 assert_eq!(record.origin, NamespaceOrigin::WorkerMint);
2380
2381 let listed = store.list_namespaces().await?;
2383 assert!(
2384 listed.iter().any(|record| record.name == "orders"),
2385 "list_namespaces returns the minted namespace"
2386 );
2387
2388 Ok(())
2389 }
2390
2391 #[cfg(feature = "haematite-backend")]
2392 #[tokio::test(flavor = "multi_thread")]
2393 async fn connect_store_haematite_round_trips_through_event_store()
2394 -> Result<(), Box<dyn std::error::Error>> {
2395 use aion_core::{ContentType, EventEnvelope, PackageVersion, Payload, RunId, WorkflowId};
2396 use aion_store::WriteToken;
2397 use chrono::Utc;
2398
2399 use crate::config::{StoreBackend, StoreConfig};
2400
2401 let data_dir = crate::test_support::private_tempdir()?;
2402 let connected = super::connect_store(StoreConfig {
2406 backend: StoreBackend::Haematite,
2407 url: None,
2408 owned_shards: Vec::new(),
2409 data_dir: Some(data_dir.path().to_string_lossy().into_owned()),
2410 shard_count: 1,
2411 cluster: None,
2412 })
2413 .await?;
2414 let event_store = connected.event_store;
2415 assert!(
2416 connected.outbox_store.is_some(),
2417 "the haematite backend shares its leaf store as the dispatcher's outbox store"
2418 );
2419 assert!(
2420 connected.bootstrap_coordinator,
2421 "a single-node haematite boot owns all shards and bootstraps the coordinator"
2422 );
2423 assert!(
2424 connected.cluster_responder.is_none(),
2425 "a single-node (no [cluster]) haematite boot has no distributed responder"
2426 );
2427
2428 let workflow_id = WorkflowId::new_v4();
2429 let event = aion_core::Event::WorkflowStarted {
2430 envelope: EventEnvelope {
2431 seq: 1,
2432 recorded_at: Utc::now(),
2433 workflow_id: workflow_id.clone(),
2434 },
2435 workflow_type: String::from("checkout"),
2436 input: Payload::new(ContentType::Json, b"{}".to_vec()),
2437 run_id: RunId::new_v4(),
2438 parent_run_id: None,
2439 package_version: PackageVersion::new("a".repeat(64)),
2440 };
2441 event_store
2442 .append(
2443 WriteToken::recorder(),
2444 &workflow_id,
2445 std::slice::from_ref(&event),
2446 0,
2447 )
2448 .await?;
2449 let history = event_store.read_history(&workflow_id).await?;
2450 assert_eq!(
2451 history.len(),
2452 1,
2453 "an event appended through the server's dyn EventStore reads back"
2454 );
2455 Ok(())
2456 }
2457
2458 #[cfg(all(feature = "haematite-backend", unix))]
2459 #[test]
2460 fn haematite_root_swap_before_first_backend_touch_cannot_redirect_writes()
2461 -> Result<(), Box<dyn std::error::Error>> {
2462 use std::os::unix::fs::symlink;
2463
2464 let sandbox = crate::test_support::private_tempdir()?;
2465 let configured_root = sandbox.path().join("data");
2466 let held_root = sandbox.path().join("held-data");
2467 let outside = sandbox.path().join("outside");
2468 std::fs::create_dir(&outside)?;
2469 let configured = configured_root
2470 .to_str()
2471 .ok_or("temporary data path was not UTF-8")?;
2472
2473 let (store, responder) =
2474 super::build_haematite_store_with_hook(configured, 4, None, || {
2475 std::fs::rename(&configured_root, &held_root)?;
2480 symlink(&outside, &configured_root)?;
2481 Ok(())
2482 })?;
2483 assert!(responder.is_none());
2484
2485 let outside_entries = std::fs::read_dir(&outside)?.collect::<Result<Vec<_>, _>>()?;
2486 assert!(
2487 outside_entries.is_empty(),
2488 "Haematite followed the replaced ambient root and wrote outside"
2489 );
2490 assert!(held_root.join("config.json").is_file());
2491 for shard in 0..4 {
2492 let shard_path = held_root.join(format!("shard-{shard}"));
2493 assert!(shard_path.is_dir(), "shard {shard} was not materialized");
2494 assert!(
2495 std::fs::read_dir(&shard_path)?
2496 .next()
2497 .transpose()?
2498 .is_some(),
2499 "shard {shard} did not run Haematite's materialization path"
2500 );
2501 }
2502
2503 drop(store);
2504 Ok(())
2505 }
2506
2507 #[cfg(all(
2508 feature = "haematite-backend",
2509 any(target_os = "linux", target_os = "android")
2510 ))]
2511 #[tokio::test]
2512 async fn proc_fd_backend_path_survives_a_post_startup_root_swap()
2513 -> Result<(), Box<dyn std::error::Error>> {
2514 use std::os::unix::fs::symlink;
2515
2516 use aion_core::{ContentType, EventEnvelope, PackageVersion, Payload, RunId, WorkflowId};
2517 use aion_store::{WritableEventStore as _, WriteToken};
2518 use chrono::Utc;
2519
2520 let sandbox = crate::test_support::private_tempdir()?;
2521 let configured_root = sandbox.path().join("data");
2522 let held_root = sandbox.path().join("held-data");
2523 let capture = sandbox.path().join("capture");
2524 std::fs::create_dir(&capture)?;
2525 let configured = configured_root
2526 .to_str()
2527 .ok_or("temporary data path was not UTF-8")?;
2528
2529 let (store, responder) = super::build_haematite_store(configured, 4, None)?;
2530 assert!(responder.is_none());
2531 std::fs::rename(&configured_root, &held_root)?;
2532 symlink(&capture, &configured_root)?;
2533
2534 let workflow_id = WorkflowId::new_v4();
2535 let event = aion_core::Event::WorkflowStarted {
2536 envelope: EventEnvelope {
2537 seq: 1,
2538 recorded_at: Utc::now(),
2539 workflow_id: workflow_id.clone(),
2540 },
2541 workflow_type: String::from("post-startup-root-swap"),
2542 input: Payload::new(ContentType::Json, b"{}".to_vec()),
2543 run_id: RunId::new_v4(),
2544 parent_run_id: None,
2545 package_version: PackageVersion::new("a".repeat(64)),
2546 };
2547 store
2548 .append(
2549 WriteToken::recorder(),
2550 &workflow_id,
2551 std::slice::from_ref(&event),
2552 0,
2553 )
2554 .await?;
2555
2556 let captured = std::fs::read_dir(&capture)?.collect::<Result<Vec<_>, _>>()?;
2557 assert!(
2558 captured.is_empty(),
2559 "post-startup append followed the replacement symlink into capture"
2560 );
2561 assert!(held_root.join("config.json").is_file());
2562 drop(store);
2563 Ok(())
2564 }
2565
2566 #[cfg(all(
2567 feature = "haematite-backend",
2568 unix,
2569 not(any(target_os = "linux", target_os = "android"))
2570 ))]
2571 #[test]
2572 fn path_ambient_haematite_refuses_group_or_world_writable_ancestors()
2573 -> Result<(), Box<dyn std::error::Error>> {
2574 use std::os::unix::fs::PermissionsExt as _;
2575
2576 let sandbox = crate::test_support::private_tempdir()?;
2577 std::fs::set_permissions(sandbox.path(), std::fs::Permissions::from_mode(0o700))?;
2578
2579 for mode in [0o770, 0o1777] {
2580 let shared = sandbox.path().join(format!("shared-{mode:o}"));
2581 let data_root = shared.join("data");
2582 std::fs::create_dir(&shared)?;
2583 std::fs::set_permissions(&shared, std::fs::Permissions::from_mode(mode))?;
2584 std::fs::create_dir(&data_root)?;
2585 std::fs::set_permissions(&data_root, std::fs::Permissions::from_mode(0o700))?;
2586 let configured = data_root
2587 .to_str()
2588 .ok_or("temporary data path was not UTF-8")?;
2589
2590 let Err(error) = super::build_haematite_store(configured, 4, None) else {
2591 return Err(format!("mode {mode:04o} ancestor was accepted").into());
2592 };
2593 let message = error.to_string();
2594 let crate::ServerError::UnsafeDataRootAncestor {
2595 data_root: resolved_root,
2596 component,
2597 reason,
2598 } = error
2599 else {
2600 return Err(format!("expected typed unsafe-ancestor error, got {message}").into());
2601 };
2602 assert_eq!(resolved_root, std::fs::canonicalize(&data_root)?);
2603 assert_eq!(component, std::fs::canonicalize(&shared)?);
2604 assert!(
2605 reason.contains(&format!("mode {mode:04o}")),
2606 "unexpected reason: {reason}"
2607 );
2608 if mode & 0o1000 != 0 {
2609 assert!(reason.contains("sticky bit is not accepted"));
2610 }
2611 assert!(message.contains("private Aion home"));
2612 assert!(
2613 !data_root.join("config.json").exists(),
2614 "Haematite touched its ambient path before the refusal"
2615 );
2616 }
2617 Ok(())
2618 }
2619
2620 #[cfg(all(feature = "haematite-backend", target_os = "macos"))]
2621 #[test]
2622 fn path_ambient_haematite_refuses_mutating_allow_acl_ancestor()
2623 -> Result<(), Box<dyn std::error::Error>> {
2624 use std::os::unix::fs::PermissionsExt as _;
2625
2626 let sandbox = crate::test_support::private_tempdir()?;
2627 std::fs::set_permissions(sandbox.path(), std::fs::Permissions::from_mode(0o700))?;
2628 let shared = sandbox.path().join("acl-shared");
2629 let data_root = shared.join("data");
2630 std::fs::create_dir(&shared)?;
2631 std::fs::set_permissions(&shared, std::fs::Permissions::from_mode(0o700))?;
2632 let acl = "everyone allow list,search,add_file,add_subdirectory,delete_child";
2633 let status = std::process::Command::new("chmod")
2634 .arg("+a")
2635 .arg(acl)
2636 .arg(&shared)
2637 .status()?;
2638 assert!(status.success(), "failed to install Darwin regression ACL");
2639 let configured = data_root
2640 .to_str()
2641 .ok_or("temporary data path was not UTF-8")?;
2642
2643 let result = super::build_haematite_store(configured, 4, None);
2644 let cleanup = std::process::Command::new("chmod")
2645 .arg("-RN")
2646 .arg(&shared)
2647 .status()?;
2648 assert!(cleanup.success(), "failed to clean Darwin regression ACL");
2649
2650 let Err(error) = result else {
2651 return Err("mutating non-euid allow ACL ancestor was accepted".into());
2652 };
2653 let message = error.to_string();
2654 let crate::ServerError::UnsafeDataRootAncestor {
2655 component, reason, ..
2656 } = error
2657 else {
2658 return Err(format!("expected typed unsafe-ancestor error, got {message}").into());
2659 };
2660 assert_eq!(component, std::fs::canonicalize(&shared)?);
2661 assert!(
2662 reason.contains("allow"),
2663 "reason did not name the ACE: {reason}"
2664 );
2665 assert!(
2666 reason.contains("everyone"),
2667 "reason did not name the ACE principal: {reason}"
2668 );
2669 assert!(
2670 !data_root.join("config.json").exists(),
2671 "Haematite touched its ambient path before the ACL refusal"
2672 );
2673 Ok(())
2674 }
2675
2676 #[cfg(all(feature = "haematite-backend", target_os = "macos"))]
2677 #[test]
2678 fn path_ambient_haematite_accepts_the_euid_uuid_allow_ace()
2679 -> Result<(), Box<dyn std::error::Error>> {
2680 use std::os::unix::fs::PermissionsExt as _;
2681
2682 use exacl::{AclEntry, AclOption, Perm};
2683
2684 let sandbox = crate::test_support::private_tempdir()?;
2685 std::fs::set_permissions(sandbox.path(), std::fs::Permissions::from_mode(0o700))?;
2686 let private_parent = sandbox.path().join("euid-uuid-allow");
2687 let data_root = private_parent.join("data");
2688 std::fs::create_dir(&private_parent)?;
2689 std::fs::set_permissions(&private_parent, std::fs::Permissions::from_mode(0o700))?;
2690
2691 let server_uid = rustix::process::geteuid().as_raw();
2692 let ace_qualifier = crate::filesystem::darwin_user_uuid_for_test(server_uid)?;
2693 let entry = AclEntry::allow_user(
2694 &ace_qualifier.to_string(),
2695 Perm::EXECUTE | Perm::WRITE | Perm::APPEND | Perm::DELETE_CHILD,
2696 None,
2697 );
2698 exacl::setfacl(
2699 &[private_parent.as_path()],
2700 &[entry],
2701 AclOption::SYMLINK_ACL,
2702 )?;
2703 let configured = data_root
2704 .to_str()
2705 .ok_or("temporary data path was not UTF-8")?;
2706
2707 let result = super::build_haematite_store(configured, 4, None);
2708 let cleanup = std::process::Command::new("chmod")
2709 .arg("-RN")
2710 .arg(&private_parent)
2711 .status()?;
2712 assert!(cleanup.success(), "failed to clean euid UUID allow ACL");
2713
2714 let (store, responder) = result?;
2715 assert!(responder.is_none());
2716 assert!(data_root.join("config.json").is_file());
2717 drop(store);
2718 Ok(())
2719 }
2720
2721 #[cfg(all(feature = "haematite-backend", target_os = "macos"))]
2722 #[test]
2723 fn path_ambient_haematite_refuses_a_non_euid_user_uuid_allow_ace()
2724 -> Result<(), Box<dyn std::error::Error>> {
2725 use std::os::unix::fs::PermissionsExt as _;
2726
2727 use exacl::{AclEntry, AclOption, Perm};
2728
2729 let sandbox = crate::test_support::private_tempdir()?;
2730 std::fs::set_permissions(sandbox.path(), std::fs::Permissions::from_mode(0o700))?;
2731 let shared = sandbox.path().join("non-euid-uuid-allow");
2732 let data_root = shared.join("data");
2733 std::fs::create_dir(&shared)?;
2734 std::fs::set_permissions(&shared, std::fs::Permissions::from_mode(0o700))?;
2735
2736 let server_uid = rustix::process::geteuid().as_raw();
2737 let foreign_uid = u32::from(server_uid == 0);
2738 let foreign_qualifier = crate::filesystem::darwin_user_uuid_for_test(foreign_uid)?;
2739 let entry = AclEntry::allow_user(
2740 &foreign_qualifier.to_string(),
2741 Perm::EXECUTE | Perm::WRITE | Perm::APPEND | Perm::DELETE_CHILD,
2742 None,
2743 );
2744 exacl::setfacl(&[shared.as_path()], &[entry], AclOption::SYMLINK_ACL)?;
2745 let configured = data_root
2746 .to_str()
2747 .ok_or("temporary data path was not UTF-8")?;
2748
2749 let result = super::build_haematite_store(configured, 4, None);
2750 let cleanup = std::process::Command::new("chmod")
2751 .arg("-RN")
2752 .arg(&shared)
2753 .status()?;
2754 assert!(cleanup.success(), "failed to clean non-euid UUID allow ACL");
2755
2756 let Err(error) = result else {
2757 return Err("mutating non-euid user UUID allow ACE was accepted".into());
2758 };
2759 let message = error.to_string();
2760 let crate::ServerError::UnsafeDataRootAncestor {
2761 component, reason, ..
2762 } = error
2763 else {
2764 return Err(format!("expected typed unsafe-ancestor error, got {message}").into());
2765 };
2766 assert_eq!(component, std::fs::canonicalize(&shared)?);
2767 assert!(
2768 reason.contains("allow") && reason.contains(&format!("server euid {server_uid}")),
2769 "reason did not name the rejected ACE: {reason}"
2770 );
2771 assert!(
2772 !data_root.join("config.json").exists(),
2773 "Haematite touched its ambient path before the UUID ACL refusal"
2774 );
2775 Ok(())
2776 }
2777
2778 #[cfg(all(feature = "haematite-backend", target_os = "macos"))]
2779 #[test]
2780 fn path_ambient_haematite_accepts_a_deny_only_acl_ancestor()
2781 -> Result<(), Box<dyn std::error::Error>> {
2782 use std::os::unix::fs::PermissionsExt as _;
2783
2784 let sandbox = crate::test_support::private_tempdir()?;
2785 std::fs::set_permissions(sandbox.path(), std::fs::Permissions::from_mode(0o700))?;
2786 let private_parent = sandbox.path().join("deny-only");
2787 let data_root = private_parent.join("data");
2788 std::fs::create_dir(&private_parent)?;
2789 std::fs::set_permissions(&private_parent, std::fs::Permissions::from_mode(0o700))?;
2790 let status = std::process::Command::new("chmod")
2791 .arg("+a")
2792 .arg("everyone deny delete")
2793 .arg(&private_parent)
2794 .status()?;
2795 assert!(status.success(), "failed to install Darwin deny-only ACL");
2796 let configured = data_root
2797 .to_str()
2798 .ok_or("temporary data path was not UTF-8")?;
2799
2800 let result = super::build_haematite_store(configured, 4, None);
2801 let cleanup = std::process::Command::new("chmod")
2802 .arg("-RN")
2803 .arg(&private_parent)
2804 .status()?;
2805 assert!(cleanup.success(), "failed to clean Darwin deny-only ACL");
2806
2807 let (store, responder) = result?;
2808 assert!(responder.is_none());
2809 assert!(data_root.join("config.json").is_file());
2810 drop(store);
2811 Ok(())
2812 }
2813
2814 #[cfg(all(feature = "haematite-backend", target_os = "macos"))]
2815 #[test]
2816 fn path_ambient_haematite_accepts_the_stock_home_acl_chain()
2817 -> Result<(), Box<dyn std::error::Error>> {
2818 use std::os::unix::fs::PermissionsExt as _;
2819 use users::os::unix::UserExt as _;
2820
2821 let effective_uid = rustix::process::geteuid().as_raw();
2822 let effective_user = users::get_user_by_uid(effective_uid)
2823 .ok_or_else(|| format!("server euid {effective_uid} has no account record"))?;
2824 let sandbox = tempfile::Builder::new()
2825 .prefix(".aion-acl-home-proof-")
2826 .tempdir_in(effective_user.home_dir())?;
2827 std::fs::set_permissions(sandbox.path(), std::fs::Permissions::from_mode(0o700))?;
2828 let data_root = sandbox.path().join("data");
2829 let configured = data_root
2830 .to_str()
2831 .ok_or("temporary data path was not UTF-8")?;
2832
2833 let (store, responder) = super::build_haematite_store(configured, 4, None)?;
2834 assert!(responder.is_none());
2835 assert!(data_root.join("config.json").is_file());
2836 drop(store);
2837 Ok(())
2838 }
2839
2840 #[cfg(all(
2841 feature = "haematite-backend",
2842 unix,
2843 not(any(target_os = "linux", target_os = "android"))
2844 ))]
2845 #[test]
2846 fn path_ambient_haematite_accepts_an_owner_controlled_chain()
2847 -> Result<(), Box<dyn std::error::Error>> {
2848 use std::os::unix::fs::PermissionsExt as _;
2849
2850 let sandbox = crate::test_support::private_tempdir()?;
2851 std::fs::set_permissions(sandbox.path(), std::fs::Permissions::from_mode(0o700))?;
2852 let private_parent = sandbox.path().join("private");
2853 let data_root = private_parent.join("data");
2854 std::fs::create_dir(&private_parent)?;
2855 std::fs::set_permissions(&private_parent, std::fs::Permissions::from_mode(0o700))?;
2856 let configured = data_root
2857 .to_str()
2858 .ok_or("temporary data path was not UTF-8")?;
2859
2860 let (store, responder) = super::build_haematite_store(configured, 4, None)?;
2861 assert!(responder.is_none());
2862 assert!(data_root.join("config.json").is_file());
2863 for shard in 0..4 {
2864 assert!(data_root.join(format!("shard-{shard}")).is_dir());
2865 }
2866 drop(store);
2867 Ok(())
2868 }
2869
2870 #[tokio::test]
2871 async fn connect_store_memory_backend_exposes_no_outbox_store()
2872 -> Result<(), Box<dyn std::error::Error>> {
2873 use crate::config::{StoreBackend, StoreConfig};
2874
2875 let connected = super::connect_store(StoreConfig {
2878 backend: StoreBackend::Memory,
2879 url: None,
2880 owned_shards: Vec::new(),
2881 data_dir: None,
2882 shard_count: 1,
2883 cluster: None,
2884 })
2885 .await?;
2886 assert!(
2887 connected.outbox_store.is_none(),
2888 "the in-memory backend exposes no outbox store"
2889 );
2890 Ok(())
2891 }
2892
2893 #[cfg(feature = "libsql-backend")]
2897 #[tokio::test]
2898 async fn connect_store_shares_outbox_store_only_for_libsql()
2899 -> Result<(), Box<dyn std::error::Error>> {
2900 use crate::config::{StoreBackend, StoreConfig};
2901
2902 let path = std::env::temp_dir().join(format!(
2907 "aion-connect-store-{}-{}.db",
2908 std::process::id(),
2909 std::time::SystemTime::now()
2910 .duration_since(std::time::UNIX_EPOCH)
2911 .map(|elapsed| elapsed.as_nanos())
2912 .unwrap_or_default()
2913 ));
2914 let connected = super::connect_store(StoreConfig {
2915 backend: StoreBackend::LibSql,
2916 url: Some(path.to_string_lossy().into_owned()),
2917 owned_shards: Vec::new(),
2918 data_dir: None,
2919 shard_count: 1,
2920 cluster: None,
2921 })
2922 .await?;
2923 assert!(
2924 connected.outbox_store.is_some(),
2925 "the libSQL backend shares its leaf store as the dispatcher's outbox store"
2926 );
2927 Ok(())
2928 }
2929
2930 #[tokio::test]
2931 async fn state_build_fails_without_event_broadcast_capacity()
2932 -> Result<(), Box<dyn std::error::Error>> {
2933 let mut runtime = runtime_config();
2934 runtime.websocket.event_broadcast_capacity = None;
2935
2936 let error = ServerState::build_with_store(InMemoryStore::default(), runtime)
2937 .await
2938 .err()
2939 .ok_or("state build must fail when event streaming is unsized")?;
2940
2941 assert!(error.is_config(), "expected a config error, got {error}");
2942 assert!(
2943 error
2944 .to_string()
2945 .contains("websocket.event_broadcast_capacity"),
2946 "error must name the missing key: {error}"
2947 );
2948 Ok(())
2949 }
2950
2951 #[tokio::test]
2952 async fn state_build_fails_without_query_timeout() -> Result<(), Box<dyn std::error::Error>> {
2953 let mut runtime = runtime_config();
2954 runtime.query_timeout = None;
2955
2956 let error = ServerState::build_with_store(InMemoryStore::default(), runtime)
2957 .await
2958 .err()
2959 .ok_or("state build must fail when the query reply deadline is unset")?;
2960
2961 assert!(error.is_config(), "expected a config error, got {error}");
2962 assert!(
2963 error.to_string().contains("runtime.query_timeout_ms"),
2964 "error must name the missing key: {error}"
2965 );
2966 assert!(
2967 error.to_string().contains("AION_RUNTIME_QUERY_TIMEOUT_MS"),
2968 "error must name the environment override: {error}"
2969 );
2970 Ok(())
2971 }
2972
2973 #[tokio::test]
2974 async fn state_build_fails_with_zero_query_timeout() -> Result<(), Box<dyn std::error::Error>> {
2975 let mut runtime = runtime_config();
2976 runtime.query_timeout = Some(Duration::ZERO);
2977
2978 let error = ServerState::build_with_store(InMemoryStore::default(), runtime)
2979 .await
2980 .err()
2981 .ok_or("state build must fail when the query reply deadline is zero")?;
2982
2983 assert!(error.is_config(), "expected a config error, got {error}");
2984 assert!(
2985 error.to_string().contains("runtime.query_timeout_ms"),
2986 "error must name the zero-valued key: {error}"
2987 );
2988 Ok(())
2989 }
2990}