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
10use crate::dev_ui::{ActivityMockRegistry, DevMockingDispatcher};
11
12#[cfg(feature = "auth")]
13use crate::auth::JwksCache;
14use crate::{
15 config::{RuntimeConfig, ServerConfig, StoreBackend, StoreConfig},
16 error::ServerError,
17 namespace::{NamespaceGuard, NamespaceMinter, resolver::NamespaceResolver},
18 observability::{
19 Metrics, health::HealthState, instrumented_store::InstrumentedEventStore,
20 metrics::MetricsError,
21 },
22 shutdown::DrainState,
23 worker::{
24 ConnectedWorkerRegistry, HeartbeatTracker, PendingActivities, WorkerActivityDispatcher,
25 supervisor::WorkerSupervisor,
26 },
27};
28
29fn new_supervisor(
40 store: &Arc<dyn WorkerDeploymentStore>,
41 publisher: &crate::cluster_publisher::ClusterEventPublisher,
42) -> Arc<WorkerSupervisor> {
43 Arc::new(WorkerSupervisor::new(Arc::clone(store), publisher.clone()))
44}
45
46#[derive(Clone)]
48pub struct ServerState {
49 inner: Arc<ServerStateInner>,
50}
51
52struct ServerStateInner {
53 namespace_guard: NamespaceGuard,
54 runtime: RuntimeConfig,
55 worker_registry: ConnectedWorkerRegistry,
56 pending_activities: PendingActivities,
57 heartbeat_tracker: HeartbeatTracker,
58 grpc_liveness_waiters: crate::worker::GrpcLivenessWaiters,
65 drain_state: DrainState,
66 metrics: Option<Metrics>,
67 health: Option<HealthState>,
68 activity_mock_registry: Option<ActivityMockRegistry>,
71 outbox_store: Option<Arc<dyn OutboxStore>>,
74 namespace_store: Arc<dyn NamespaceStore>,
83 worker_deployment_store: Arc<dyn WorkerDeploymentStore>,
85 worker_supervisor: Arc<WorkerSupervisor>,
90 outbox_wake: Arc<tokio::sync::Notify>,
97 cluster_publisher: crate::cluster_publisher::ClusterEventPublisher,
101 transcript_publisher: crate::activity_publisher::ActivityEventPublisher,
110 attempt_owners: crate::worker::AttemptOwnerIndex,
117 queue_service_state: crate::worker::QueueServiceState,
123 queue_declarations: crate::worker::QueueDeclarationSource,
127 update_status: crate::update_check::UpdateStatusState,
133 workspace_root: crate::worker::WorkspaceRoot,
141 cluster_self_node: Option<String>,
146 cluster_responder: Option<aion_store_haematite::ClusterResponder>,
151 cluster_store: Option<Arc<aion_store_haematite::HaematiteStore>>,
154 watched_peers: Vec<crate::cluster::WatchedPeer>,
157 shard_directory: Option<Arc<crate::routing::StaticShardDirectory>>,
162 request_forwarder: Option<Arc<dyn crate::routing::RequestForwarder>>,
166 #[cfg(feature = "auth")]
167 jwks_cache: Option<JwksCache>,
168}
169
170impl ServerState {
171 const FALLBACK_CLUSTER_BROADCAST_CAPACITY: std::num::NonZeroUsize =
180 match std::num::NonZeroUsize::new(64) {
181 Some(value) => value,
182 None => std::num::NonZeroUsize::MIN,
183 };
184
185 pub async fn build(config: ServerConfig) -> Result<Self, ServerError> {
192 let (store_config, runtime) = config.into_parts();
193 let connected = connect_store(store_config).await?;
194 Self::build_with_connected_store(connected, runtime).await
195 }
196
197 pub async fn build_with_store<S>(store: S, runtime: RuntimeConfig) -> Result<Self, ServerError>
203 where
204 S: EventStore + NamespaceStore + WorkerDeploymentStore,
205 {
206 let leaf = Arc::new(store);
210 let namespace_store: Arc<dyn NamespaceStore> = leaf.clone();
211 let worker_deployment_store: Arc<dyn WorkerDeploymentStore> = leaf.clone();
212 Self::build_with_connected_store(
213 ConnectedStore::local(leaf, None, namespace_store, worker_deployment_store),
214 runtime,
215 )
216 .await
217 }
218
219 async fn build_with_connected_store(
220 connected: ConnectedStore,
221 runtime: RuntimeConfig,
222 ) -> Result<Self, ServerError> {
223 let cluster_self_node = connected.cluster_self_node();
224 let outbox_store = connected.outbox_store;
225 let bootstrap_coordinator = connected.bootstrap_coordinator;
226 let cluster_responder = connected.cluster_responder;
227 let cluster_store = connected.cluster_store;
228 let watched_peers = connected.watched_peers;
229 let RoutingState {
233 shard_directory,
234 request_forwarder,
235 mint_routing,
236 } = build_routing_state(
237 cluster_store.as_ref(),
238 connected.directory_peers,
239 connected.self_node_id,
240 );
241 let (event_broadcast_capacity, query_timeout) = required_engine_seams(&runtime)?;
242 let (cluster_publisher, transcript_publisher) =
243 build_real_time_publishers(&runtime, connected.observability_store)?;
244 let (metrics, outbox_wake, instrumented_store) =
245 build_instrumented_store(&runtime, connected.event_store)?;
246 let exported_metrics = runtime.metrics.enabled.then_some(metrics.clone());
247 let seams = build_worker_seams(
248 &runtime,
249 &cluster_publisher,
250 &connected.namespace_store,
251 &connected.worker_deployment_store,
252 mint_routing,
253 );
254 let (
255 activity_dispatcher,
256 activity_mock_registry,
257 attempt_owners,
258 workspace_root,
259 update_status,
260 ) = build_decorated_dispatcher(&runtime, &seams, transcript_publisher.clone());
261
262 let engine = boot_engine(EngineAssembly {
263 seams: &seams,
264 instrumented_store: &instrumented_store,
265 event_broadcast_capacity,
266 query_timeout,
267 activity_dispatcher,
268 active_registry: Arc::new(aion::Registry::default()),
269 bootstrap_coordinator,
270 runtime: &runtime,
271 })
272 .await?;
273 let resolver = NamespaceResolver::from_config(runtime.namespace.clone(), engine);
274 let worker_supervisor =
275 new_supervisor(&connected.worker_deployment_store, &cluster_publisher);
276 #[cfg(feature = "auth")]
277 let jwks_cache = build_jwks_cache(&runtime).await?;
278 Ok(Self {
279 inner: Arc::new(ServerStateInner {
280 namespace_guard: NamespaceGuard::new(resolver),
281 runtime,
282 metrics: exported_metrics,
283 worker_registry: seams.worker_registry,
284 pending_activities: seams.pending_activities,
285 heartbeat_tracker: seams.heartbeat_tracker,
286 grpc_liveness_waiters: crate::worker::GrpcLivenessWaiters::new(),
287 drain_state: seams.drain_state,
288 health: Some(HealthState::new(instrumented_store, true)),
289 activity_mock_registry,
290 outbox_store,
291 namespace_store: connected.namespace_store,
292 worker_supervisor,
293 worker_deployment_store: connected.worker_deployment_store,
294 outbox_wake,
295 cluster_publisher,
296 transcript_publisher,
297 attempt_owners,
298 queue_service_state: seams.queue_service_state,
299 queue_declarations: seams.queue_declarations,
300 update_status,
301 workspace_root,
302 cluster_self_node,
303 cluster_responder,
304 cluster_store,
305 watched_peers,
306 shard_directory,
307 request_forwarder,
308 #[cfg(feature = "auth")]
309 jwks_cache,
310 }),
311 })
312 }
313
314 #[must_use]
316 pub fn from_parts(namespace_resolver: NamespaceResolver, runtime: RuntimeConfig) -> Self {
317 Self::from_parts_with_namespace_store(
322 namespace_resolver,
323 runtime,
324 Arc::new(aion_store::InMemoryStore::default()),
325 )
326 }
327
328 #[must_use]
335 pub fn from_parts_with_namespace_store<S>(
336 namespace_resolver: NamespaceResolver,
337 runtime: RuntimeConfig,
338 store: Arc<S>,
339 ) -> Self
340 where
341 S: NamespaceStore + WorkerDeploymentStore,
342 {
343 let namespace_store: Arc<dyn NamespaceStore> = store.clone();
344 let worker_deployment_store: Arc<dyn WorkerDeploymentStore> = store;
345 Self::from_parts_with_control_stores(
346 namespace_resolver,
347 runtime,
348 namespace_store,
349 worker_deployment_store,
350 )
351 }
352
353 #[must_use]
358 pub fn from_parts_with_control_stores(
359 namespace_resolver: NamespaceResolver,
360 runtime: RuntimeConfig,
361 namespace_store: Arc<dyn NamespaceStore>,
362 worker_deployment_store: Arc<dyn WorkerDeploymentStore>,
363 ) -> Self {
364 let heartbeat_tracker = HeartbeatTracker::new(runtime.worker.heartbeat_window);
365 let pending_activities =
372 PendingActivities::default().with_heartbeat_window(runtime.worker.heartbeat_window);
373 let bounds = transcript_bounds(&runtime);
377 let batch = required_transcript_batch_policy(&runtime)
381 .unwrap_or(crate::activity_publisher::TranscriptBatchPolicy::UNBATCHED);
382 let cluster_publisher = crate::cluster_publisher::ClusterEventPublisher::new(
383 Self::FALLBACK_CLUSTER_BROADCAST_CAPACITY,
384 );
385 Self {
386 inner: Arc::new(ServerStateInner {
387 namespace_guard: NamespaceGuard::new(namespace_resolver),
388 runtime,
389 worker_registry: ConnectedWorkerRegistry::default()
390 .with_worker_deployment_store(worker_deployment_store.clone())
391 .with_cluster_publisher(cluster_publisher.clone()),
392 pending_activities,
393 heartbeat_tracker,
394 grpc_liveness_waiters: crate::worker::GrpcLivenessWaiters::new(),
395 drain_state: DrainState::default(),
396 metrics: None,
397 health: None,
398 activity_mock_registry: None,
399 outbox_store: None,
400 namespace_store,
401 worker_supervisor: new_supervisor(&worker_deployment_store, &cluster_publisher),
402 worker_deployment_store,
403 outbox_wake: Arc::new(tokio::sync::Notify::new()),
404 cluster_publisher,
405 transcript_publisher: build_transcript_publisher(
409 None,
410 Self::FALLBACK_CLUSTER_BROADCAST_CAPACITY,
411 bounds,
412 batch,
413 ),
414 attempt_owners: crate::worker::AttemptOwnerIndex::new(),
415 queue_service_state: crate::worker::QueueServiceState::default(),
416 queue_declarations: crate::worker::QueueDeclarationSource::default(),
417 update_status: crate::update_check::UpdateStatusState::default(),
421 workspace_root: crate::worker::WorkspaceRoot::resolve(),
425 cluster_self_node: None,
426 cluster_responder: None,
427 cluster_store: None,
428 watched_peers: Vec::new(),
429 shard_directory: None,
430 request_forwarder: None,
431 #[cfg(feature = "auth")]
432 jwks_cache: None,
433 }),
434 }
435 }
436
437 #[cfg(feature = "auth")]
446 #[must_use]
447 pub fn from_parts_with_namespace_store_and_jwks<S>(
448 namespace_resolver: NamespaceResolver,
449 runtime: RuntimeConfig,
450 store: Arc<S>,
451 jwks_cache: JwksCache,
452 ) -> Self
453 where
454 S: NamespaceStore + WorkerDeploymentStore,
455 {
456 let namespace_store: Arc<dyn NamespaceStore> = store.clone();
457 let worker_deployment_store: Arc<dyn WorkerDeploymentStore> = store;
458 let heartbeat_tracker = HeartbeatTracker::new(runtime.worker.heartbeat_window);
459 let pending_activities =
466 PendingActivities::default().with_heartbeat_window(runtime.worker.heartbeat_window);
467 let bounds = transcript_bounds(&runtime);
471 let batch = required_transcript_batch_policy(&runtime)
475 .unwrap_or(crate::activity_publisher::TranscriptBatchPolicy::UNBATCHED);
476 let cluster_publisher = crate::cluster_publisher::ClusterEventPublisher::new(
477 Self::FALLBACK_CLUSTER_BROADCAST_CAPACITY,
478 );
479 Self {
480 inner: Arc::new(ServerStateInner {
481 namespace_guard: NamespaceGuard::new(namespace_resolver),
482 runtime,
483 worker_registry: ConnectedWorkerRegistry::default()
484 .with_worker_deployment_store(worker_deployment_store.clone())
485 .with_cluster_publisher(cluster_publisher.clone()),
486 pending_activities,
487 heartbeat_tracker,
488 grpc_liveness_waiters: crate::worker::GrpcLivenessWaiters::new(),
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 batch,
507 ),
508 attempt_owners: crate::worker::AttemptOwnerIndex::new(),
509 queue_service_state: crate::worker::QueueServiceState::default(),
510 queue_declarations: crate::worker::QueueDeclarationSource::default(),
511 update_status: crate::update_check::UpdateStatusState::default(),
515 workspace_root: crate::worker::WorkspaceRoot::resolve(),
519 cluster_self_node: None,
520 cluster_responder: None,
521 cluster_store: None,
522 watched_peers: Vec::new(),
523 shard_directory: None,
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 batch = required_transcript_batch_policy(&runtime)
564 .unwrap_or(crate::activity_publisher::TranscriptBatchPolicy::UNBATCHED);
565 let cluster_publisher = crate::cluster_publisher::ClusterEventPublisher::new(
566 Self::FALLBACK_CLUSTER_BROADCAST_CAPACITY,
567 );
568 Self {
569 inner: Arc::new(ServerStateInner {
570 namespace_guard: NamespaceGuard::new(namespace_resolver),
571 runtime,
572 worker_registry: ConnectedWorkerRegistry::default(),
573 pending_activities,
574 heartbeat_tracker,
575 grpc_liveness_waiters: crate::worker::GrpcLivenessWaiters::new(),
576 drain_state: DrainState::default(),
577 metrics: None,
578 health: None,
579 activity_mock_registry: None,
580 outbox_store: None,
581 namespace_store: Arc::new(aion_store::InMemoryStore::default()),
586 worker_supervisor: new_supervisor(&fallback_deployment_store, &cluster_publisher),
587 worker_deployment_store: fallback_deployment_store,
588 outbox_wake: Arc::new(tokio::sync::Notify::new()),
589 cluster_publisher,
590 transcript_publisher: build_transcript_publisher(
594 None,
595 Self::FALLBACK_CLUSTER_BROADCAST_CAPACITY,
596 bounds,
597 batch,
598 ),
599 attempt_owners: crate::worker::AttemptOwnerIndex::new(),
600 queue_service_state: crate::worker::QueueServiceState::default(),
601 queue_declarations: crate::worker::QueueDeclarationSource::default(),
602 update_status: crate::update_check::UpdateStatusState::default(),
606 workspace_root: crate::worker::WorkspaceRoot::resolve(),
610 cluster_self_node: None,
611 cluster_responder: None,
612 cluster_store: None,
613 watched_peers: Vec::new(),
614 shard_directory: None,
615 request_forwarder: None,
616 jwks_cache: Some(jwks_cache),
617 }),
618 }
619 }
620
621 #[must_use]
623 pub fn from_parts_with_registry(
624 namespace_resolver: NamespaceResolver,
625 runtime: RuntimeConfig,
626 worker_registry: ConnectedWorkerRegistry,
627 ) -> Self {
628 let fallback_deployment_store: Arc<dyn WorkerDeploymentStore> =
632 Arc::new(aion_store::InMemoryStore::default());
633 let heartbeat_tracker = HeartbeatTracker::new(runtime.worker.heartbeat_window);
634 let pending_activities =
641 PendingActivities::default().with_heartbeat_window(runtime.worker.heartbeat_window);
642 let bounds = transcript_bounds(&runtime);
646 let batch = required_transcript_batch_policy(&runtime)
650 .unwrap_or(crate::activity_publisher::TranscriptBatchPolicy::UNBATCHED);
651 let cluster_publisher = crate::cluster_publisher::ClusterEventPublisher::new(
652 Self::FALLBACK_CLUSTER_BROADCAST_CAPACITY,
653 );
654 Self {
655 inner: Arc::new(ServerStateInner {
656 namespace_guard: NamespaceGuard::new(namespace_resolver),
657 runtime,
658 worker_registry,
659 pending_activities,
660 heartbeat_tracker,
661 grpc_liveness_waiters: crate::worker::GrpcLivenessWaiters::new(),
662 drain_state: DrainState::default(),
663 metrics: None,
664 health: None,
665 activity_mock_registry: None,
666 outbox_store: None,
667 namespace_store: Arc::new(aion_store::InMemoryStore::default()),
672 worker_supervisor: new_supervisor(&fallback_deployment_store, &cluster_publisher),
673 worker_deployment_store: fallback_deployment_store,
674 outbox_wake: Arc::new(tokio::sync::Notify::new()),
675 cluster_publisher,
676 transcript_publisher: build_transcript_publisher(
680 None,
681 Self::FALLBACK_CLUSTER_BROADCAST_CAPACITY,
682 bounds,
683 batch,
684 ),
685 attempt_owners: crate::worker::AttemptOwnerIndex::new(),
686 queue_service_state: crate::worker::QueueServiceState::default(),
687 queue_declarations: crate::worker::QueueDeclarationSource::default(),
688 update_status: crate::update_check::UpdateStatusState::default(),
692 workspace_root: crate::worker::WorkspaceRoot::resolve(),
696 cluster_self_node: None,
697 cluster_responder: None,
698 cluster_store: None,
699 watched_peers: Vec::new(),
700 shard_directory: None,
701 request_forwarder: None,
702 #[cfg(feature = "auth")]
703 jwks_cache: None,
704 }),
705 }
706 }
707
708 #[must_use]
710 pub fn namespace_guard(&self) -> &NamespaceGuard {
711 &self.inner.namespace_guard
712 }
713
714 #[must_use]
716 pub fn deploy_guard(&self) -> crate::deploy::DeployGuard {
717 crate::deploy::DeployGuard::new(self.inner.namespace_guard.resolver().clone())
718 }
719
720 #[must_use]
722 pub fn runtime_config(&self) -> &RuntimeConfig {
723 &self.inner.runtime
724 }
725
726 #[must_use]
734 pub(crate) fn update_status(&self) -> &crate::update_check::UpdateStatusState {
735 &self.inner.update_status
736 }
737
738 #[must_use]
748 pub fn workspace_root(&self) -> &crate::worker::WorkspaceRoot {
749 &self.inner.workspace_root
750 }
751
752 #[must_use]
754 pub fn worker_registry(&self) -> &ConnectedWorkerRegistry {
755 &self.inner.worker_registry
756 }
757
758 pub fn cancel_in_flight_activities(
773 &self,
774 workflow_id: &aion_core::WorkflowId,
775 ) -> Result<Vec<crate::worker::CancelRequest>, ServerError> {
776 crate::worker::cancel_in_flight_activities(
777 self.heartbeat_tracker(),
778 self.worker_registry(),
779 workflow_id,
780 )
781 }
782
783 #[must_use]
787 pub fn cluster_publisher(&self) -> &crate::cluster_publisher::ClusterEventPublisher {
788 &self.inner.cluster_publisher
789 }
790
791 #[must_use]
796 pub fn transcript_publisher(&self) -> &crate::activity_publisher::ActivityEventPublisher {
797 &self.inner.transcript_publisher
798 }
799
800 #[must_use]
804 pub fn attempt_owners(&self) -> &crate::worker::AttemptOwnerIndex {
805 &self.inner.attempt_owners
806 }
807
808 #[must_use]
811 pub fn queue_service_state(&self) -> &crate::worker::QueueServiceState {
812 &self.inner.queue_service_state
813 }
814
815 #[must_use]
819 pub fn queue_declarations(&self) -> &crate::worker::QueueDeclarationSource {
820 &self.inner.queue_declarations
821 }
822
823 pub fn unserved_queues(&self) -> Result<Vec<crate::worker::UnservedQueue>, ServerError> {
834 self.inner.queue_service_state.unserved()
835 }
836
837 pub fn unrecoverable_runs(
855 &self,
856 ) -> Result<Vec<(aion_core::WorkflowId, aion::registry::UnrecoverableRun)>, ServerError> {
857 self.engine()?
858 .registry()
859 .unrecoverable()
860 .list()
861 .map_err(ServerError::from)
862 }
863
864 #[must_use]
876 pub fn intervention_router(&self) -> crate::worker::InterventionRouter {
877 let transport: std::sync::Arc<dyn crate::worker::InterventionTransport> = {
878 #[cfg(feature = "liminal-transport")]
879 {
880 std::sync::Arc::new(crate::worker::LiminalInterventionTransport)
881 }
882 #[cfg(not(feature = "liminal-transport"))]
883 {
884 std::sync::Arc::new(NullInterventionTransport)
885 }
886 };
887 crate::worker::InterventionRouter::new(
888 self.inner.worker_registry.clone(),
889 self.inner.attempt_owners.clone(),
890 transport,
891 )
892 .with_transcript_publisher(self.inner.transcript_publisher.clone())
895 }
896
897 #[must_use]
901 pub fn cluster_self_node(&self) -> Option<&str> {
902 self.inner.cluster_self_node.as_deref()
903 }
904
905 pub fn engine(&self) -> Result<Arc<aion::Engine>, ServerError> {
918 self.inner
919 .namespace_guard
920 .resolver()
921 .engine()
922 .map(Arc::clone)
923 }
924
925 #[must_use]
927 pub fn pending_activities(&self) -> &PendingActivities {
928 &self.inner.pending_activities
929 }
930
931 #[must_use]
933 pub fn heartbeat_tracker(&self) -> &HeartbeatTracker {
934 &self.inner.heartbeat_tracker
935 }
936
937 #[must_use]
943 pub fn grpc_liveness_waiters(&self) -> &crate::worker::GrpcLivenessWaiters {
944 &self.inner.grpc_liveness_waiters
945 }
946
947 #[must_use]
949 pub fn drain_state(&self) -> &DrainState {
950 &self.inner.drain_state
951 }
952
953 #[must_use]
955 pub fn metrics(&self) -> Option<&Metrics> {
956 self.inner.metrics.as_ref()
957 }
958
959 #[must_use]
961 pub fn health(&self) -> Option<&HealthState> {
962 self.inner.health.as_ref()
963 }
964
965 #[must_use]
970 pub fn activity_mock_registry(&self) -> Option<&ActivityMockRegistry> {
971 self.inner.activity_mock_registry.as_ref()
972 }
973
974 #[must_use]
980 pub fn outbox_store(&self) -> Option<Arc<dyn OutboxStore>> {
981 self.inner.outbox_store.clone()
982 }
983
984 #[must_use]
993 pub fn namespace_store(&self) -> &Arc<dyn NamespaceStore> {
994 &self.inner.namespace_store
995 }
996
997 #[must_use]
999 pub fn worker_deployment_store(&self) -> &Arc<dyn WorkerDeploymentStore> {
1000 &self.inner.worker_deployment_store
1001 }
1002
1003 #[must_use]
1010 pub fn worker_supervisor(&self) -> &Arc<WorkerSupervisor> {
1011 &self.inner.worker_supervisor
1012 }
1013
1014 #[must_use]
1023 pub fn namespace_minter(&self) -> NamespaceMinter {
1024 let minter = NamespaceMinter::new(
1025 Arc::clone(&self.inner.namespace_store),
1026 self.inner.runtime.auto_create,
1027 )
1028 .with_cluster_publisher(self.inner.cluster_publisher.clone());
1033 match self.namespace_routing() {
1038 Some(routing) => minter.with_routing(routing),
1039 None => minter,
1040 }
1041 }
1042
1043 #[must_use]
1053 pub fn namespace_routing(&self) -> Option<crate::namespace::NamespaceRouting> {
1054 build_namespace_routing(
1055 self.cluster_store(),
1056 self.shard_directory(),
1057 self.request_forwarder(),
1058 )
1059 }
1060
1061 #[must_use]
1067 pub fn outbox_wake(&self) -> Arc<tokio::sync::Notify> {
1068 Arc::clone(&self.inner.outbox_wake)
1069 }
1070
1071 #[must_use]
1077 pub fn is_clustered(&self) -> bool {
1078 self.inner.cluster_responder.is_some()
1079 }
1080
1081 #[must_use]
1086 pub fn cluster_store(&self) -> Option<&Arc<aion_store_haematite::HaematiteStore>> {
1087 self.inner.cluster_store.as_ref()
1088 }
1089
1090 #[must_use]
1094 pub fn shard_directory(&self) -> Option<&Arc<crate::routing::StaticShardDirectory>> {
1095 self.inner.shard_directory.as_ref()
1096 }
1097
1098 #[must_use]
1101 pub fn request_forwarder(&self) -> Option<&Arc<dyn crate::routing::RequestForwarder>> {
1102 self.inner.request_forwarder.as_ref()
1103 }
1104
1105 #[must_use]
1122 pub fn spawn_heartbeat_sweeper(
1123 &self,
1124 shutdown: tokio::sync::watch::Receiver<bool>,
1125 ) -> tokio::task::JoinHandle<()> {
1126 let sweeper = crate::worker::HeartbeatSweeper::new(
1127 self.inner.heartbeat_tracker.clone(),
1128 self.inner.worker_registry.clone(),
1129 self.inner.pending_activities.clone(),
1130 self.inner.drain_state.clone(),
1131 self.inner.runtime.worker.heartbeat_window,
1132 )
1133 .with_queue_state(self.inner.queue_service_state.clone());
1137 tokio::spawn(sweeper.run(shutdown))
1138 }
1139
1140 #[cfg(feature = "liminal-transport")]
1151 #[must_use]
1152 pub fn spawn_liminal_liveness_probe(
1153 &self,
1154 notifier: std::sync::Arc<crate::worker::LiminalConnectionNotifier>,
1155 shutdown: tokio::sync::watch::Receiver<bool>,
1156 ) -> tokio::task::JoinHandle<()> {
1157 let probe = crate::worker::LivenessProbe::across_transports(
1158 Some(notifier),
1159 Some(self.inner.grpc_liveness_waiters.clone()),
1164 self.inner.heartbeat_tracker.clone(),
1165 self.inner.worker_registry.clone(),
1166 self.inner.runtime.worker.heartbeat_window,
1167 );
1168 tokio::spawn(probe.run(shutdown))
1169 }
1170
1171 pub fn spawn_cluster_supervisor(
1187 &self,
1188 config: crate::cluster::SupervisorConfig,
1189 shutdown: tokio::sync::watch::Receiver<bool>,
1190 ) -> Result<bool, ServerError> {
1191 let Some(cluster_store) = self.inner.cluster_store.clone() else {
1192 return Ok(false);
1193 };
1194 if self.inner.watched_peers.is_empty() {
1195 return Ok(false);
1196 }
1197 let engine = Arc::clone(self.inner.namespace_guard.resolver().engine()?);
1198 let publisher = Arc::new(self.inner.cluster_publisher.clone());
1202 let self_node = self.inner.cluster_self_node.clone().unwrap_or_default();
1203 let adopter = Arc::new(crate::cluster::OutboxSettlingAdopter::new(
1209 engine,
1210 self.inner.outbox_store.clone(),
1211 ));
1212 let supervisor = crate::cluster::ClusterSupervisor::new(
1213 cluster_store,
1214 adopter,
1215 self.inner.watched_peers.clone(),
1216 config,
1217 )
1218 .with_publisher(publisher, self_node);
1219 if !supervisor.watches_any() {
1220 return Ok(false);
1221 }
1222 tokio::spawn(supervisor.run(shutdown));
1223 Ok(true)
1224 }
1225
1226 #[cfg(feature = "auth")]
1228 #[must_use]
1229 pub fn jwks_cache(&self) -> Option<&JwksCache> {
1230 self.inner.jwks_cache.as_ref()
1231 }
1232
1233 pub fn shutdown(&self) -> Result<(), ServerError> {
1240 self.inner.namespace_guard.resolver().shutdown_engine()
1241 }
1242}
1243
1244#[cfg(feature = "auth")]
1245async fn build_jwks_cache(runtime: &RuntimeConfig) -> Result<Option<JwksCache>, ServerError> {
1246 if !runtime.auth.enabled {
1247 return Ok(None);
1248 }
1249 let Some(url) = runtime.auth.jwks_url.clone() else {
1250 return Err(ServerError::Config {
1251 message: "auth.jwks_url must not be empty when auth.enabled is true".to_owned(),
1252 });
1253 };
1254 let interval = std::time::Duration::from_secs(runtime.auth.jwks_refresh_seconds);
1255 let cache = JwksCache::new(url, interval)
1256 .await
1257 .map_err(|error| ServerError::Config {
1258 message: format!("auth jwks initial fetch failed: {error}"),
1259 })?;
1260 Ok(Some(cache))
1261}
1262
1263fn metrics_config_error(error: &MetricsError) -> ServerError {
1264 ServerError::Config {
1265 message: error.to_string(),
1266 }
1267}
1268
1269struct EngineAssembly<'a> {
1271 seams: &'a WorkerSeams,
1273 instrumented_store: &'a Arc<InstrumentedEventStore>,
1275 event_broadcast_capacity: std::num::NonZeroUsize,
1277 query_timeout: std::time::Duration,
1279 activity_dispatcher: Arc<dyn ActivityDispatcher>,
1281 active_registry: Arc<aion::Registry>,
1283 bootstrap_coordinator: bool,
1285 runtime: &'a RuntimeConfig,
1287}
1288
1289fn server_search_attribute_schema() -> Result<aion_core::SearchAttributeSchema, ServerError> {
1298 let mut schema = aion_core::SearchAttributeSchema::new();
1299 for (name, label) in [
1300 (crate::namespace::NAMESPACE_ATTRIBUTE, "namespace"),
1301 (crate::namespace::TASK_QUEUE_ATTRIBUTE, "task_queue"),
1302 (crate::namespace::DISPLAY_NAME_ATTRIBUTE, "display_name"),
1303 ] {
1304 schema
1305 .register(name, aion_core::SearchAttributeType::String)
1306 .map_err(|error| ServerError::Config {
1307 message: format!("failed to register {label} search attribute: {error}"),
1308 })?;
1309 }
1310 Ok(schema)
1311}
1312
1313async fn boot_engine(assembly: EngineAssembly<'_>) -> Result<Arc<aion::Engine>, ServerError> {
1323 let search_attribute_schema = server_search_attribute_schema()?;
1324 let runtime = assembly.runtime;
1325 let builder = EngineBuilder::new()
1326 .store_arc(assembly.instrumented_store.clone())
1327 .event_streaming(assembly.event_broadcast_capacity)
1328 .in_memory_visibility()
1329 .search_attribute_schema(search_attribute_schema)
1330 .scheduler_threads(runtime.scheduler_threads)
1331 .outbox_enabled(runtime.outbox.enabled)
1332 .activity_dispatcher(assembly.activity_dispatcher)
1333 .active_registry(assembly.active_registry)
1334 .production_recovery_seam()
1335 .defer_startup_recovery()
1342 .signal_router_factory(|runtime: Arc<RuntimeHandle>, handoff| {
1343 Arc::new(ConcreteSignalRouter::new(runtime, handoff)) as Arc<dyn SignalRouter>
1344 })
1345 .query_timeout(assembly.query_timeout)
1346 .bootstrap_schedule_coordinator(assembly.bootstrap_coordinator)
1352 .load_workflow_sources(runtime.workflow_packages.iter().map(PathBuf::as_path));
1353 let builder = if runtime.owned_shards.is_empty() {
1359 builder
1360 } else {
1361 builder.owned_shards(runtime.owned_shards.iter().copied())
1362 };
1363 let engine = Arc::new(builder.build().await.map_err(ServerError::from)?);
1364 install_engine_backed_seams(assembly.seams, &engine, runtime.outbox.enabled);
1365 engine
1371 .run_startup_recovery()
1372 .await
1373 .map_err(ServerError::from)?;
1374 Ok(engine)
1375}
1376
1377fn required_engine_seams(
1382 runtime: &RuntimeConfig,
1383) -> Result<(std::num::NonZeroUsize, std::time::Duration), ServerError> {
1384 let event_broadcast_capacity = runtime
1385 .websocket
1386 .event_broadcast_capacity
1387 .and_then(std::num::NonZeroUsize::new)
1388 .ok_or_else(|| ServerError::Config {
1389 message: crate::config::EVENT_BROADCAST_CAPACITY_REQUIRED.to_owned(),
1390 })?;
1391 let query_timeout = runtime
1392 .query_timeout
1393 .filter(|timeout| !timeout.is_zero())
1394 .ok_or_else(|| ServerError::Config {
1395 message: crate::config::QUERY_TIMEOUT_REQUIRED.to_owned(),
1396 })?;
1397 Ok((event_broadcast_capacity, query_timeout))
1398}
1399
1400fn install_outbox_delivery(
1407 pending_activities: &PendingActivities,
1408 engine: &Arc<aion::Engine>,
1409 outbox_enabled: bool,
1410) {
1411 if outbox_enabled {
1412 let callback = Arc::new(crate::worker::ServerOutboxDeliveryCallback::new(
1413 Arc::clone(engine),
1414 ));
1415 pending_activities.set_outbox_delivery(callback);
1416 }
1417}
1418
1419struct WorkerSeams {
1421 worker_registry: ConnectedWorkerRegistry,
1422 pending_activities: PendingActivities,
1423 heartbeat_tracker: HeartbeatTracker,
1424 drain_state: DrainState,
1425 queue_declarations: crate::worker::QueueDeclarationSource,
1426 queue_service_state: crate::worker::QueueServiceState,
1427 declared_bodies: crate::worker::DeclaredBodySource,
1428 cluster_publisher: crate::cluster_publisher::ClusterEventPublisher,
1432}
1433
1434fn build_worker_seams(
1439 runtime: &RuntimeConfig,
1440 cluster_publisher: &crate::cluster_publisher::ClusterEventPublisher,
1441 namespace_store: &Arc<dyn NamespaceStore>,
1442 worker_deployment_store: &Arc<dyn WorkerDeploymentStore>,
1443 mint_routing: Option<crate::namespace::NamespaceRouting>,
1444) -> WorkerSeams {
1445 let worker_registry = ConnectedWorkerRegistry::default()
1446 .with_cluster_publisher(cluster_publisher.clone())
1447 .with_namespace_minting(namespace_store.clone(), runtime.auto_create)
1448 .with_worker_deployment_store(worker_deployment_store.clone());
1449 let worker_registry = match mint_routing {
1455 Some(routing) => worker_registry.with_namespace_routing(routing),
1456 None => worker_registry,
1457 };
1458 WorkerSeams {
1459 worker_registry,
1460 pending_activities: PendingActivities::default()
1465 .with_heartbeat_window(runtime.worker.heartbeat_window),
1466 heartbeat_tracker: HeartbeatTracker::new(runtime.worker.heartbeat_window),
1467 drain_state: DrainState::default(),
1468 queue_declarations: crate::worker::QueueDeclarationSource::default(),
1469 queue_service_state: crate::worker::QueueServiceState::default(),
1470 declared_bodies: crate::worker::DeclaredBodySource::default(),
1471 cluster_publisher: cluster_publisher.clone(),
1472 }
1473}
1474
1475fn install_queue_declarations(
1482 queue_declarations: &crate::worker::QueueDeclarationSource,
1483 engine: &Arc<aion::Engine>,
1484) {
1485 queue_declarations.install(Arc::new(crate::worker::EngineQueueDeclarations::new(
1486 Arc::clone(engine),
1487 )));
1488}
1489
1490fn install_declared_bodies(
1505 declared_bodies: &crate::worker::DeclaredBodySource,
1506 engine: &Arc<aion::Engine>,
1507) {
1508 declared_bodies.install(Arc::new(crate::worker::EngineDeclaredBodies::new(
1509 Arc::clone(engine),
1510 )));
1511}
1512
1513fn install_engine_backed_seams(seams: &WorkerSeams, engine: &Arc<aion::Engine>, outbox: bool) {
1519 install_outbox_delivery(&seams.pending_activities, engine, outbox);
1520 install_queue_declarations(&seams.queue_declarations, engine);
1521 install_declared_bodies(&seams.declared_bodies, engine);
1522}
1523
1524fn build_bridge_dispatcher(
1537 runtime: &RuntimeConfig,
1538 seams: &WorkerSeams,
1539) -> (WorkerActivityDispatcher, crate::worker::AttemptOwnerIndex) {
1540 let attempt_owners = crate::worker::AttemptOwnerIndex::new();
1541 let dispatcher = WorkerActivityDispatcher::new(
1542 seams.worker_registry.clone(),
1543 runtime.default_namespace.clone(),
1544 seams.heartbeat_tracker.clone(),
1545 )
1546 .with_pending(seams.pending_activities.clone())
1547 .with_drain_state(seams.drain_state.clone())
1548 .with_tokio_handle(tokio::runtime::Handle::current())
1549 .with_attempt_owners(attempt_owners.clone())
1550 .with_queue_service(runtime.worker.queue_service.clone())
1551 .with_queue_declarations(seams.queue_declarations.clone())
1552 .with_queue_state(seams.queue_service_state.clone())
1553 .with_cluster_publisher(seams.cluster_publisher.clone());
1554 (dispatcher, attempt_owners)
1555}
1556
1557fn build_instrumented_store(
1570 runtime: &RuntimeConfig,
1571 event_store: Arc<dyn EventStore>,
1572) -> Result<
1573 (
1574 Metrics,
1575 Arc<tokio::sync::Notify>,
1576 Arc<InstrumentedEventStore>,
1577 ),
1578 ServerError,
1579> {
1580 let metrics = Metrics::new().map_err(|error| metrics_config_error(&error))?;
1581 let outbox_wake = Arc::new(tokio::sync::Notify::new());
1582 let instrumented_store = Arc::new(
1583 InstrumentedEventStore::new(
1584 event_store,
1585 metrics.clone(),
1586 runtime.default_namespace.clone(),
1587 )
1588 .with_outbox_wake(Arc::clone(&outbox_wake)),
1589 );
1590 Ok((metrics, outbox_wake, instrumented_store))
1591}
1592
1593fn build_decorated_dispatcher(
1598 runtime: &RuntimeConfig,
1599 seams: &WorkerSeams,
1600 transcript: crate::activity_publisher::ActivityEventPublisher,
1601) -> (
1602 Arc<dyn ActivityDispatcher>,
1603 Option<ActivityMockRegistry>,
1604 crate::worker::AttemptOwnerIndex,
1605 crate::worker::WorkspaceRoot,
1606 crate::update_check::UpdateStatusState,
1607) {
1608 let workspace_root = crate::worker::WorkspaceRoot::resolve();
1612 let update_status = crate::update_check::UpdateStatusState::default();
1615 let (dispatcher, attempt_owners) = build_bridge_dispatcher(runtime, seams);
1616 let (activity_dispatcher, activity_mock_registry) = decorate_activity_dispatcher(
1617 dispatcher,
1618 seams.declared_bodies.clone(),
1619 workspace_root.clone(),
1620 transcript,
1621 update_status.clone(),
1622 runtime.dev.enabled,
1623 );
1624 (
1625 activity_dispatcher,
1626 activity_mock_registry,
1627 attempt_owners,
1628 workspace_root,
1629 update_status,
1630 )
1631}
1632
1633fn decorate_activity_dispatcher(
1634 dispatcher: WorkerActivityDispatcher,
1635 declared_bodies: crate::worker::DeclaredBodySource,
1636 workspace_root: crate::worker::WorkspaceRoot,
1637 transcript: crate::activity_publisher::ActivityEventPublisher,
1638 update_status: crate::update_check::UpdateStatusState,
1639 dev_enabled: bool,
1640) -> (Arc<dyn ActivityDispatcher>, Option<ActivityMockRegistry>) {
1641 let declared = crate::worker::DeclaredCommandDispatcher::new(
1649 Arc::new(dispatcher),
1650 declared_bodies.clone(),
1651 tokio::runtime::Handle::current(),
1652 workspace_root,
1653 transcript,
1654 );
1655 let observed = crate::update_check::UpdateCheckObserver::new(
1656 Arc::new(declared),
1657 declared_bodies,
1658 update_status,
1659 );
1660 if dev_enabled {
1661 let registry = ActivityMockRegistry::new();
1662 let decorated = DevMockingDispatcher::new(Arc::new(observed), registry.clone());
1663 (Arc::new(decorated), Some(registry))
1664 } else {
1665 (Arc::new(observed), None)
1666 }
1667}
1668
1669fn required_cluster_broadcast_capacity(
1674 runtime: &RuntimeConfig,
1675) -> Result<std::num::NonZeroUsize, ServerError> {
1676 runtime
1677 .websocket
1678 .cluster_broadcast_capacity
1679 .and_then(std::num::NonZeroUsize::new)
1680 .ok_or_else(|| ServerError::Config {
1681 message: crate::config::CLUSTER_BROADCAST_CAPACITY_REQUIRED.to_owned(),
1682 })
1683}
1684
1685fn build_real_time_publishers(
1698 runtime: &RuntimeConfig,
1699 observability_store: Option<Arc<dyn aion_store::ObservabilityStore>>,
1700) -> Result<
1701 (
1702 crate::cluster_publisher::ClusterEventPublisher,
1703 crate::activity_publisher::ActivityEventPublisher,
1704 ),
1705 ServerError,
1706> {
1707 let capacity = required_cluster_broadcast_capacity(runtime)?;
1708 let batch = required_transcript_batch_policy(runtime)?;
1709 Ok((
1710 crate::cluster_publisher::ClusterEventPublisher::new(capacity),
1711 build_transcript_publisher(
1712 observability_store,
1713 capacity,
1714 transcript_bounds(runtime),
1715 batch,
1716 ),
1717 ))
1718}
1719
1720fn required_node_cache_budget(
1750 config: &StoreConfig,
1751) -> Result<haematite::NodeCacheBudget, ServerError> {
1752 config.node_cache_budget.ok_or_else(|| ServerError::Config {
1753 message: crate::config::STORE_NODE_CACHE_BUDGET_REQUIRED.to_owned(),
1754 })
1755}
1756
1757fn required_transcript_batch_policy(
1758 runtime: &RuntimeConfig,
1759) -> Result<crate::activity_publisher::TranscriptBatchPolicy, ServerError> {
1760 let max_batch_events = runtime
1761 .observability
1762 .max_batch_events
1763 .and_then(std::num::NonZeroUsize::new)
1764 .ok_or_else(|| ServerError::Config {
1765 message: crate::config::OBSERVABILITY_MAX_BATCH_EVENTS_REQUIRED.to_owned(),
1766 })?;
1767 let max_batch_hold_ms =
1770 runtime
1771 .observability
1772 .max_batch_hold_ms
1773 .ok_or_else(|| ServerError::Config {
1774 message: crate::config::OBSERVABILITY_MAX_BATCH_HOLD_MS_REQUIRED.to_owned(),
1775 })?;
1776 Ok(crate::activity_publisher::TranscriptBatchPolicy {
1777 max_batch_events,
1778 max_hold: std::time::Duration::from_millis(max_batch_hold_ms),
1779 })
1780}
1781
1782fn transcript_bounds(runtime: &RuntimeConfig) -> crate::activity_bounds::TranscriptBounds {
1784 crate::activity_bounds::TranscriptBounds {
1785 max_event_bytes: runtime.observability.max_event_bytes,
1786 max_stream_events: runtime.observability.max_stream_events,
1787 }
1788}
1789
1790fn build_transcript_publisher(
1801 observability_store: Option<Arc<dyn aion_store::ObservabilityStore>>,
1802 capacity: std::num::NonZeroUsize,
1803 bounds: crate::activity_bounds::TranscriptBounds,
1804 batch: crate::activity_publisher::TranscriptBatchPolicy,
1805) -> crate::activity_publisher::ActivityEventPublisher {
1806 let store = observability_store
1807 .unwrap_or_else(|| Arc::new(aion_store::InMemoryObservabilityStore::default()));
1808 crate::activity_publisher::ActivityEventPublisher::new(store, capacity, batch)
1809 .with_bounds(bounds)
1810}
1811
1812struct RoutingState {
1814 shard_directory: Option<Arc<crate::routing::StaticShardDirectory>>,
1815 request_forwarder: Option<Arc<dyn crate::routing::RequestForwarder>>,
1816 mint_routing: Option<crate::namespace::NamespaceRouting>,
1821}
1822
1823fn build_routing_state(
1827 cluster_store: Option<&Arc<aion_store_haematite::HaematiteStore>>,
1828 directory_peers: Vec<crate::routing::DirectoryPeer>,
1829 self_node_id: Option<String>,
1830) -> RoutingState {
1831 let Some(store) = cluster_store else {
1832 return RoutingState {
1833 shard_directory: None,
1834 request_forwarder: None,
1835 mint_routing: None,
1836 };
1837 };
1838 let shard_directory = Arc::new(crate::routing::StaticShardDirectory::new(
1839 Arc::clone(store),
1840 directory_peers,
1841 self_node_id,
1842 ));
1843 let request_forwarder: Arc<dyn crate::routing::RequestForwarder> =
1844 Arc::new(crate::routing::GrpcRequestForwarder::new());
1845 let mint_routing = build_namespace_routing(
1846 Some(store),
1847 Some(&shard_directory),
1848 Some(&request_forwarder),
1849 );
1850 RoutingState {
1851 shard_directory: Some(shard_directory),
1852 request_forwarder: Some(request_forwarder),
1853 mint_routing,
1854 }
1855}
1856
1857fn build_namespace_routing(
1867 cluster_store: Option<&Arc<aion_store_haematite::HaematiteStore>>,
1868 shard_directory: Option<&Arc<crate::routing::StaticShardDirectory>>,
1869 request_forwarder: Option<&Arc<dyn crate::routing::RequestForwarder>>,
1870) -> Option<crate::namespace::NamespaceRouting> {
1871 use crate::namespace::{
1872 GrpcMintForwarder, MintForwarder, MintShardOwners, NamespaceRouting, NamespaceShardResolver,
1873 };
1874 let store = Arc::clone(cluster_store?);
1875 let directory = Arc::clone(shard_directory?);
1876 let shards: Arc<dyn NamespaceShardResolver> = store;
1877 let owners: Arc<dyn MintShardOwners> = directory;
1878 let forwarder: Arc<dyn MintForwarder> =
1879 Arc::new(GrpcMintForwarder::new(Arc::clone(request_forwarder?)));
1880 Some(NamespaceRouting::new(shards, owners, forwarder))
1881}
1882
1883struct ConnectedStore {
1893 event_store: Arc<dyn EventStore>,
1894 outbox_store: Option<Arc<dyn OutboxStore>>,
1895 namespace_store: Arc<dyn NamespaceStore>,
1902 worker_deployment_store: Arc<dyn WorkerDeploymentStore>,
1904 observability_store: Option<Arc<dyn aion_store::ObservabilityStore>>,
1912 bootstrap_coordinator: bool,
1913 cluster_responder: Option<aion_store_haematite::ClusterResponder>,
1914 cluster_store: Option<Arc<aion_store_haematite::HaematiteStore>>,
1918 watched_peers: Vec<crate::cluster::WatchedPeer>,
1921 directory_peers: Vec<crate::routing::DirectoryPeer>,
1925 self_node_id: Option<String>,
1929}
1930
1931impl ConnectedStore {
1932 fn local(
1939 event_store: Arc<dyn EventStore>,
1940 outbox_store: Option<Arc<dyn OutboxStore>>,
1941 namespace_store: Arc<dyn NamespaceStore>,
1942 worker_deployment_store: Arc<dyn WorkerDeploymentStore>,
1943 ) -> Self {
1944 Self {
1945 event_store,
1946 outbox_store,
1947 namespace_store,
1948 worker_deployment_store,
1949 observability_store: None,
1955 bootstrap_coordinator: true,
1956 cluster_responder: None,
1957 cluster_store: None,
1958 watched_peers: Vec::new(),
1959 directory_peers: Vec::new(),
1960 self_node_id: None,
1961 }
1962 }
1963
1964 fn cluster_self_node(&self) -> Option<String> {
1967 self.self_node_id.clone()
1968 }
1969}
1970
1971async fn connect_store(config: StoreConfig) -> Result<ConnectedStore, ServerError> {
1974 match config.backend {
1975 StoreBackend::Memory => {
1976 let leaf = Arc::new(aion_store::InMemoryStore::default());
1979 let namespace_store: Arc<dyn NamespaceStore> = leaf.clone();
1980 let worker_deployment_store: Arc<dyn WorkerDeploymentStore> = leaf.clone();
1981 Ok(ConnectedStore::local(
1982 leaf,
1983 None,
1984 namespace_store,
1985 worker_deployment_store,
1986 ))
1987 }
1988 StoreBackend::Haematite => connect_haematite_store(config).await,
1989 }
1990}
1991
1992async fn connect_haematite_store(config: StoreConfig) -> Result<ConnectedStore, ServerError> {
2012 let node_cache_budget = required_node_cache_budget(&config)?;
2018 let Some(data_dir) = config.data_dir else {
2019 return Err(ServerError::Config {
2020 message: "store.data_dir must not be empty when store.backend is haematite".to_owned(),
2021 });
2022 };
2023 let shard_count = config.shard_count;
2024 let owned_shards = config.owned_shards.clone();
2025 let cluster = config.cluster.clone();
2026 let watched_peers: Vec<crate::cluster::WatchedPeer> = cluster
2031 .as_ref()
2032 .map(|cluster| {
2033 cluster
2034 .peers
2035 .iter()
2036 .map(|peer| crate::cluster::WatchedPeer {
2037 name: peer.name.clone(),
2038 owned_shards: peer.owned_shards.clone(),
2039 })
2040 .collect()
2041 })
2042 .unwrap_or_default();
2043 let directory_peers: Vec<crate::routing::DirectoryPeer> = cluster
2046 .as_ref()
2047 .map(|cluster| {
2048 cluster
2049 .peers
2050 .iter()
2051 .map(|peer| crate::routing::DirectoryPeer {
2052 name: peer.name.clone(),
2053 owned_shards: peer.owned_shards.clone(),
2054 grpc_addr: peer.grpc_address,
2055 })
2056 .collect()
2057 })
2058 .unwrap_or_default();
2059 let self_node_id: Option<String> = cluster.as_ref().map(|cluster| cluster.node_id.clone());
2062 let (store, responder) = tokio::task::spawn_blocking(move || {
2066 build_haematite_store(&data_dir, shard_count, cluster, node_cache_budget)
2067 })
2068 .await
2069 .map_err(|error| ServerError::Config {
2070 message: format!("haematite store initialization task failed: {error}"),
2071 })??;
2072
2073 let bootstrap_coordinator = if owned_shards.is_empty() {
2077 true
2078 } else {
2079 store.set_owned_shards(owned_shards.iter().copied());
2080 store.owns_workflow_shard(&aion::schedule_coordinator_workflow_id())
2081 };
2082
2083 let leaf = Arc::new(store);
2084 let event_store: Arc<dyn EventStore> = leaf.clone();
2085 let outbox_store: Arc<dyn OutboxStore> = leaf.clone();
2086 let namespace_store: Arc<dyn NamespaceStore> = leaf.clone();
2090 let worker_deployment_store: Arc<dyn WorkerDeploymentStore> = leaf.clone();
2091 let observability_store: Arc<dyn aion_store::ObservabilityStore> = leaf.clone();
2096 let cluster_store = responder.as_ref().map(|_| leaf);
2100 let (watched_peers, directory_peers, self_node_id) = if cluster_store.is_some() {
2101 (watched_peers, directory_peers, self_node_id)
2102 } else {
2103 (Vec::new(), Vec::new(), None)
2104 };
2105 Ok(ConnectedStore {
2106 event_store,
2107 outbox_store: Some(outbox_store),
2108 namespace_store,
2109 worker_deployment_store,
2110 observability_store: Some(observability_store),
2111 bootstrap_coordinator,
2112 cluster_responder: responder,
2113 cluster_store,
2114 watched_peers,
2115 directory_peers,
2116 self_node_id,
2117 })
2118}
2119
2120fn build_haematite_store(
2135 data_dir: &str,
2136 shard_count: usize,
2137 cluster: Option<crate::config::ClusterConfig>,
2138 node_cache_budget: haematite::NodeCacheBudget,
2139) -> Result<
2140 (
2141 aion_store_haematite::HaematiteStore,
2142 Option<aion_store_haematite::ClusterResponder>,
2143 ),
2144 ServerError,
2145> {
2146 build_haematite_store_with_hook(data_dir, shard_count, cluster, node_cache_budget, || Ok(()))
2147}
2148
2149fn build_haematite_store_with_hook(
2150 data_dir: &str,
2151 shard_count: usize,
2152 cluster: Option<crate::config::ClusterConfig>,
2153 node_cache_budget: haematite::NodeCacheBudget,
2154 before_backend_touch: impl FnOnce() -> Result<(), std::io::Error>,
2155) -> Result<
2156 (
2157 aion_store_haematite::HaematiteStore,
2158 Option<aion_store_haematite::ClusterResponder>,
2159 ),
2160 ServerError,
2161> {
2162 use aion_store_haematite::{ClusterBootstrap, HaematiteStore};
2163
2164 let private_root = crate::filesystem::ConfinedDir::open_or_create(std::path::Path::new(
2172 data_dir,
2173 ))
2174 .map_err(|error| ServerError::Config {
2175 message: format!("unsafe store.data_dir `{data_dir}`: {error}"),
2176 })?;
2177
2178 for shard in 0..shard_count {
2182 private_root
2183 .create_dir_all(std::path::Path::new(&format!("shard-{shard}")))
2184 .map_err(|error| ServerError::Config {
2185 message: format!(
2186 "failed to materialize shard-{shard} under store.data_dir `{data_dir}`: {error}"
2187 ),
2188 })?;
2189 }
2190 private_root
2191 .harden_tree()
2192 .map_err(|error| private_store_mode_error(data_dir, &error))?;
2193
2194 before_backend_touch().map_err(|error| ServerError::Config {
2197 message: format!("store.data_dir pre-open hook failed: {error}"),
2198 })?;
2199
2200 #[cfg(unix)]
2201 let backend_path = private_root
2202 .backend_path()
2203 .map_err(|error| ServerError::Config {
2204 message: format!("failed to resolve held store.data_dir `{data_dir}`: {error}"),
2205 })?;
2206 #[cfg(all(unix, not(any(target_os = "linux", target_os = "android"))))]
2207 crate::filesystem::validate_ambient_backend_ancestors(&backend_path).map_err(|error| {
2208 let (component, reason) = error.into_parts();
2209 ServerError::UnsafeDataRootAncestor {
2210 data_root: backend_path.clone(),
2211 component,
2212 reason,
2213 }
2214 })?;
2215 #[cfg(not(unix))]
2216 let backend_path = std::path::PathBuf::from(data_dir);
2217
2218 let Some(cluster) = cluster else {
2219 let store = if backend_path.join("config.json").exists() {
2220 HaematiteStore::open(&backend_path, node_cache_budget).map_err(ServerError::from)?
2225 } else {
2226 HaematiteStore::create_with_shard_count(&backend_path, shard_count, node_cache_budget)
2227 .map_err(ServerError::from)?
2228 };
2229 store.materialize_all_shards().map_err(ServerError::from)?;
2230 private_root
2231 .harden_tree()
2232 .map_err(|error| private_store_mode_error(data_dir, &error))?;
2233 let store = store.retain_data_root_capability(private_root);
2234 return Ok((store, None));
2235 };
2236
2237 let boot = ClusterBootstrap {
2238 node_id: cluster.node_id,
2239 bind_address: cluster.bind_address,
2240 members: cluster.members,
2241 peers: cluster
2242 .peers
2243 .into_iter()
2244 .map(|peer| (peer.name, peer.address))
2245 .collect(),
2246 timeout: HAEMATITE_CLUSTER_OP_TIMEOUT,
2247 };
2248 let (store, responder) = HaematiteStore::open_or_create_distributed(
2249 &backend_path,
2250 shard_count,
2251 boot,
2252 node_cache_budget,
2253 )
2254 .map_err(ServerError::from)?;
2255 store.materialize_all_shards().map_err(ServerError::from)?;
2256 private_root
2257 .harden_tree()
2258 .map_err(|error| private_store_mode_error(data_dir, &error))?;
2259 let store = store.retain_data_root_capability(private_root);
2260 Ok((store, Some(responder)))
2261}
2262
2263fn private_store_mode_error(data_dir: &str, error: &std::io::Error) -> ServerError {
2264 ServerError::Config {
2265 message: format!(
2266 "failed to apply private modes under store.data_dir `{data_dir}`: {error}"
2267 ),
2268 }
2269}
2270
2271const HAEMATITE_CLUSTER_OP_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(5);
2273
2274#[cfg(not(feature = "liminal-transport"))]
2283#[derive(Clone, Debug)]
2284struct NullInterventionTransport;
2285
2286#[cfg(not(feature = "liminal-transport"))]
2287#[async_trait::async_trait]
2288impl crate::worker::InterventionTransport for NullInterventionTransport {
2289 async fn push(
2290 &self,
2291 _worker: &crate::worker::WorkerHandle,
2292 _command: aion_core::InterventionCommand,
2293 ) -> Result<aion_core::InterventionOutcome, ServerError> {
2294 Err(ServerError::worker_connection_lost(
2295 "intervention",
2296 "no intervention push transport is compiled in".to_owned(),
2297 ))
2298 }
2299}
2300
2301#[cfg(test)]
2302mod tests {
2303 use std::{net::SocketAddr, time::Duration};
2304
2305 use aion_store::InMemoryStore;
2306
2307 use super::ServerState;
2308 use crate::config::{
2309 AuthConfig, AuthoringConfig, DeployConfig, DevConfig, ListenConfig, MetricsConfig,
2310 NamespaceConfig, NamespaceMode, OpsConsoleAssetSource, OpsConsoleConfig, OutboxConfig,
2311 RuntimeConfig, WebSocketConfig, WorkerConfig,
2312 };
2313
2314 fn runtime_config() -> RuntimeConfig {
2315 RuntimeConfig {
2316 listen: ListenConfig {
2317 grpc: SocketAddr::from(([127, 0, 0, 1], 50051)),
2318 http: SocketAddr::from(([127, 0, 0, 1], 8080)),
2319 },
2320 tls: None,
2321 auth: AuthConfig {
2322 enabled: false,
2323 jwks_url: None,
2324 jwks_refresh_seconds: 300,
2325 },
2326 ops_console: OpsConsoleConfig {
2327 source: OpsConsoleAssetSource::Embedded,
2328 },
2329 namespace: NamespaceConfig {
2330 mode: NamespaceMode::SharedEngine,
2331 },
2332 worker: WorkerConfig {
2333 heartbeat_window: Duration::from_secs(30),
2334 ..WorkerConfig::default()
2335 },
2336 websocket: WebSocketConfig {
2337 outbound_buffer_bound: 32,
2338 event_broadcast_capacity: Some(64),
2339 cluster_broadcast_capacity: Some(64),
2340 },
2341 workflow_packages: Vec::new(),
2342 deploy: DeployConfig::default(),
2343 authoring: AuthoringConfig::default(),
2344 dev: DevConfig::default(),
2345 outbox: OutboxConfig::default(),
2346 observability: crate::config::ObservabilityConfig::with_flush_policy(64, 0),
2347 mcp: crate::config::ResolvedMcpConfig::default(),
2348 scheduler_threads: 1,
2349 query_timeout: Some(Duration::from_secs(10)),
2350 default_namespace: "default".to_owned(),
2351 auto_create: crate::config::AutoCreate::Open,
2352 max_in_flight_activities: crate::config::DEFAULT_MAX_IN_FLIGHT_ACTIVITIES,
2353 drain_timeout: Duration::from_secs(30),
2354 metrics: MetricsConfig { enabled: true },
2355 owned_shards: Vec::new(),
2356 cors_allowed_origins: Vec::new(),
2357 }
2358 }
2359
2360 #[test]
2365 fn an_unruled_flush_policy_refuses_to_build_the_publisher() {
2366 for (mutate, expected_key) in [
2367 (
2368 Box::new(|runtime: &mut RuntimeConfig| {
2369 runtime.observability.max_batch_events = None;
2370 }) as Box<dyn Fn(&mut RuntimeConfig)>,
2371 "observability.max_batch_events",
2372 ),
2373 (
2374 Box::new(|runtime: &mut RuntimeConfig| {
2375 runtime.observability.max_batch_events = Some(0);
2376 }),
2377 "observability.max_batch_events",
2378 ),
2379 (
2380 Box::new(|runtime: &mut RuntimeConfig| {
2381 runtime.observability.max_batch_hold_ms = None;
2382 }),
2383 "observability.max_batch_hold_ms",
2384 ),
2385 ] {
2386 let mut runtime = runtime_config();
2387 mutate(&mut runtime);
2388 let message = super::required_transcript_batch_policy(&runtime)
2389 .err()
2390 .map_or_else(String::new, |error| error.to_string());
2391 assert!(
2392 message.contains(expected_key),
2393 "the refusal must name {expected_key}: {message}"
2394 );
2395 assert!(
2396 message.contains("no default"),
2397 "and must say the key has no default: {message}"
2398 );
2399 }
2400 }
2401
2402 #[test]
2405 fn a_stated_flush_policy_is_carried_through_verbatim() -> Result<(), Box<dyn std::error::Error>>
2406 {
2407 let mut runtime = runtime_config();
2408 runtime.observability = crate::config::ObservabilityConfig::with_flush_policy(32, 0);
2409 let policy = super::required_transcript_batch_policy(&runtime)?;
2410 assert_eq!(policy.max_batch_events.get(), 32);
2411 assert_eq!(policy.max_hold, Duration::ZERO);
2412
2413 runtime.observability = crate::config::ObservabilityConfig::with_flush_policy(8, 250);
2414 let policy = super::required_transcript_batch_policy(&runtime)?;
2415 assert_eq!(policy.max_batch_events.get(), 8);
2416 assert_eq!(policy.max_hold, Duration::from_millis(250));
2417 Ok(())
2418 }
2419
2420 #[test]
2431 fn engine_schema_accepts_every_attribute_the_start_writer_records()
2432 -> Result<(), Box<dyn std::error::Error>> {
2433 let schema = super::server_search_attribute_schema()?;
2434 let recorded = crate::api::handlers::workflows::start_search_attributes(
2435 "tenant-a",
2436 Some("gpu"),
2437 Some("Nightly settlement"),
2438 );
2439
2440 assert!(
2441 recorded.contains_key(crate::namespace::DISPLAY_NAME_ATTRIBUTE),
2442 "the fixture must exercise the display-name attribute, or this test \
2443 cannot see its registration go missing"
2444 );
2445 for (name, value) in &recorded {
2446 schema.validate(name, value).map_err(|error| {
2447 format!(
2448 "the start writer records {name}, but the engine's schema refuses it: {error}"
2449 )
2450 })?;
2451 }
2452 Ok(())
2453 }
2454
2455 #[tokio::test]
2456 async fn builds_state_with_in_memory_store() -> Result<(), Box<dyn std::error::Error>> {
2457 let state =
2458 ServerState::build_with_store(InMemoryStore::default(), runtime_config()).await?;
2459
2460 std::hint::black_box(state.namespace_guard());
2461 std::hint::black_box(state.worker_registry());
2462
2463 Ok(())
2464 }
2465
2466 #[tokio::test]
2474 async fn unserved_queues_surfaces_a_parked_dispatch_and_clears_it()
2475 -> Result<(), Box<dyn std::error::Error>> {
2476 use aion::{ActivityDispatch, ActivityDispatcher as _};
2477 use aion_core::{ActivityId, RunId, WorkflowId};
2478 use std::sync::Arc;
2479
2480 let state =
2481 ServerState::build_with_store(InMemoryStore::default(), runtime_config()).await?;
2482 assert!(
2483 state.unserved_queues()?.is_empty(),
2484 "a calm boot has no unserved queues"
2485 );
2486 assert!(state.queue_declarations().is_installed());
2490 assert_eq!(
2491 state
2492 .queue_declarations()
2493 .declaration_for("nobody-serves-this"),
2494 crate::worker::QueueDeclaration::Unknown
2495 );
2496
2497 let dispatcher = Arc::new(
2498 crate::worker::WorkerActivityDispatcher::new(
2499 state.worker_registry().clone(),
2500 "default",
2501 crate::worker::HeartbeatTracker::new(Duration::from_secs(5)),
2502 )
2503 .with_queue_state(state.queue_service_state().clone())
2504 .with_queue_declarations(state.queue_declarations().clone()),
2505 );
2506 let workflow_id = WorkflowId::new_v4();
2507 let request = ActivityDispatch {
2508 namespace: "default".to_owned(),
2509 task_queue: "nobody-serves-this".to_owned(),
2510 node: None,
2511 workflow_id: workflow_id.clone(),
2512 run_id: RunId::new_v4(),
2513 activity_id: ActivityId::from_sequence_position(0),
2514 name: "greet".to_owned(),
2515 input: "{}".to_owned(),
2516 config: "{}".to_owned(),
2517 attempt: 1,
2518 advisory: false,
2519 labels: std::collections::BTreeMap::new(),
2520 };
2521 let parked = std::thread::spawn(move || dispatcher.dispatch(request));
2522
2523 let mut unserved = Vec::new();
2524 for _ in 0..30 {
2525 unserved = state.unserved_queues()?;
2526 if !unserved.is_empty() {
2527 break;
2528 }
2529 tokio::time::sleep(Duration::from_millis(100)).await;
2530 }
2531 assert_eq!(unserved.len(), 1, "the parked dispatch is not surfaced");
2532 assert_eq!(
2533 unserved[0].reason,
2534 crate::worker::QueueServiceReason::NoLivePollers,
2535 "an empty catalog must not be read as a structural refusal"
2536 );
2537 assert_eq!(unserved[0].key.task_queue, "nobody-serves-this");
2538 assert_eq!(unserved[0].waiting.len(), 1);
2539 assert_eq!(unserved[0].waiting[0].workflow_id, workflow_id);
2540
2541 let (worker_tx, worker_rx) = tokio::sync::mpsc::channel(1);
2543 drop(worker_rx);
2544 let registration = state.worker_registry().register_namespaces(
2545 [String::from("default")],
2546 "nobody-serves-this",
2547 None,
2548 [String::from("greet")].iter(),
2549 worker_tx,
2550 )?;
2551 let outcome = parked.join().map_err(|_| "parked dispatch panicked")?;
2552 assert!(outcome.is_err(), "the released dispatch must resolve");
2553 assert!(
2554 state.unserved_queues()?.is_empty(),
2555 "a resolved dispatch must leave the unserved state"
2556 );
2557 registration.deregister()?;
2558 Ok(())
2559 }
2560
2561 #[tokio::test]
2562 async fn namespace_store_is_reachable_and_functional_after_default_boot()
2563 -> Result<(), Box<dyn std::error::Error>> {
2564 use aion_store::{MintOutcome, NamespaceOrigin};
2565
2566 let state =
2570 ServerState::build_with_store(InMemoryStore::default(), runtime_config()).await?;
2571
2572 let store = state.namespace_store();
2573
2574 let outcome = store
2576 .register_namespace("orders", NamespaceOrigin::WorkerMint)
2577 .await?;
2578 assert_eq!(
2579 outcome,
2580 MintOutcome::Created,
2581 "the first reference to a namespace mints it"
2582 );
2583
2584 let again = store
2586 .register_namespace("orders", NamespaceOrigin::WorkerMint)
2587 .await?;
2588 assert_eq!(
2589 again,
2590 MintOutcome::AlreadyExisted,
2591 "a second reference touches the existing record rather than re-creating it"
2592 );
2593
2594 let fetched = store.get_namespace("orders").await?;
2596 let record = fetched.ok_or("registered namespace must be retrievable via get_namespace")?;
2597 assert_eq!(record.name, "orders");
2598 assert_eq!(record.origin, NamespaceOrigin::WorkerMint);
2599
2600 let listed = store.list_namespaces().await?;
2602 assert!(
2603 listed.iter().any(|record| record.name == "orders"),
2604 "list_namespaces returns the minted namespace"
2605 );
2606
2607 Ok(())
2608 }
2609
2610 #[tokio::test(flavor = "multi_thread")]
2611 async fn connect_store_haematite_round_trips_through_event_store()
2612 -> Result<(), Box<dyn std::error::Error>> {
2613 use aion_core::{ContentType, EventEnvelope, PackageVersion, Payload, RunId, WorkflowId};
2614 use aion_store::WriteToken;
2615 use chrono::Utc;
2616
2617 use crate::config::{StoreBackend, StoreConfig};
2618
2619 let data_dir = crate::test_support::private_tempdir()?;
2620 let connected = super::connect_store(StoreConfig {
2624 backend: StoreBackend::Haematite,
2625 owned_shards: Vec::new(),
2626 data_dir: Some(data_dir.path().to_string_lossy().into_owned()),
2627 shard_count: 1,
2628 cluster: None,
2629 node_cache_budget: Some(test_node_cache_budget()?),
2630 ..StoreConfig::default()
2631 })
2632 .await?;
2633 let event_store = connected.event_store;
2634 assert!(
2635 connected.outbox_store.is_some(),
2636 "the haematite backend shares its leaf store as the dispatcher's outbox store"
2637 );
2638 assert!(
2639 connected.bootstrap_coordinator,
2640 "a single-node haematite boot owns all shards and bootstraps the coordinator"
2641 );
2642 assert!(
2643 connected.cluster_responder.is_none(),
2644 "a single-node (no [cluster]) haematite boot has no distributed responder"
2645 );
2646
2647 let workflow_id = WorkflowId::new_v4();
2648 let event = aion_core::Event::WorkflowStarted {
2649 envelope: EventEnvelope {
2650 seq: 1,
2651 recorded_at: Utc::now(),
2652 workflow_id: workflow_id.clone(),
2653 },
2654 workflow_type: String::from("checkout"),
2655 input: Payload::new(ContentType::Json, b"{}".to_vec()),
2656 run_id: RunId::new_v4(),
2657 parent_run_id: None,
2658 parent_workflow_id: None,
2659 package_version: PackageVersion::new("a".repeat(64)),
2660 };
2661 event_store
2662 .append(
2663 WriteToken::recorder(),
2664 &workflow_id,
2665 std::slice::from_ref(&event),
2666 0,
2667 )
2668 .await?;
2669 let history = event_store.read_history(&workflow_id).await?;
2670 assert_eq!(
2671 history.len(),
2672 1,
2673 "an event appended through the server's dyn EventStore reads back"
2674 );
2675 Ok(())
2676 }
2677
2678 fn test_node_cache_budget() -> Result<haematite::NodeCacheBudget, Box<dyn std::error::Error>> {
2684 Ok(haematite::NodeCacheBudget::bytes(1 << 30)?)
2685 }
2686
2687 #[cfg(unix)]
2688 #[test]
2689 fn haematite_root_swap_before_first_backend_touch_cannot_redirect_writes()
2690 -> Result<(), Box<dyn std::error::Error>> {
2691 use std::os::unix::fs::symlink;
2692
2693 let sandbox = crate::test_support::private_tempdir()?;
2694 let configured_root = sandbox.path().join("data");
2695 let held_root = sandbox.path().join("held-data");
2696 let outside = sandbox.path().join("outside");
2697 std::fs::create_dir(&outside)?;
2698 let configured = configured_root
2699 .to_str()
2700 .ok_or("temporary data path was not UTF-8")?;
2701
2702 let (store, responder) = super::build_haematite_store_with_hook(
2703 configured,
2704 4,
2705 None,
2706 test_node_cache_budget()?,
2707 || {
2708 std::fs::rename(&configured_root, &held_root)?;
2713 symlink(&outside, &configured_root)?;
2714 Ok(())
2715 },
2716 )?;
2717 assert!(responder.is_none());
2718
2719 let outside_entries = std::fs::read_dir(&outside)?.collect::<Result<Vec<_>, _>>()?;
2720 assert!(
2721 outside_entries.is_empty(),
2722 "Haematite followed the replaced ambient root and wrote outside"
2723 );
2724 assert!(held_root.join("config.json").is_file());
2725 for shard in 0..4 {
2726 let shard_path = held_root.join(format!("shard-{shard}"));
2727 assert!(shard_path.is_dir(), "shard {shard} was not materialized");
2728 assert!(
2729 std::fs::read_dir(&shard_path)?
2730 .next()
2731 .transpose()?
2732 .is_some(),
2733 "shard {shard} did not run Haematite's materialization path"
2734 );
2735 }
2736
2737 drop(store);
2738 Ok(())
2739 }
2740
2741 #[cfg(any(target_os = "linux", target_os = "android"))]
2742 #[tokio::test]
2743 async fn proc_fd_backend_path_survives_a_post_startup_root_swap()
2744 -> Result<(), Box<dyn std::error::Error>> {
2745 use std::os::unix::fs::symlink;
2746
2747 use aion_core::{ContentType, EventEnvelope, PackageVersion, Payload, RunId, WorkflowId};
2748 use aion_store::{WritableEventStore as _, WriteToken};
2749 use chrono::Utc;
2750
2751 let sandbox = crate::test_support::private_tempdir()?;
2752 let configured_root = sandbox.path().join("data");
2753 let held_root = sandbox.path().join("held-data");
2754 let capture = sandbox.path().join("capture");
2755 std::fs::create_dir(&capture)?;
2756 let configured = configured_root
2757 .to_str()
2758 .ok_or("temporary data path was not UTF-8")?;
2759
2760 let (store, responder) =
2761 super::build_haematite_store(configured, 4, None, test_node_cache_budget()?)?;
2762 assert!(responder.is_none());
2763 std::fs::rename(&configured_root, &held_root)?;
2764 symlink(&capture, &configured_root)?;
2765
2766 let workflow_id = WorkflowId::new_v4();
2767 let event = aion_core::Event::WorkflowStarted {
2768 envelope: EventEnvelope {
2769 seq: 1,
2770 recorded_at: Utc::now(),
2771 workflow_id: workflow_id.clone(),
2772 },
2773 workflow_type: String::from("post-startup-root-swap"),
2774 input: Payload::new(ContentType::Json, b"{}".to_vec()),
2775 run_id: RunId::new_v4(),
2776 parent_run_id: None,
2777 package_version: PackageVersion::new("a".repeat(64)),
2778 };
2779 store
2780 .append(
2781 WriteToken::recorder(),
2782 &workflow_id,
2783 std::slice::from_ref(&event),
2784 0,
2785 )
2786 .await?;
2787
2788 let captured = std::fs::read_dir(&capture)?.collect::<Result<Vec<_>, _>>()?;
2789 assert!(
2790 captured.is_empty(),
2791 "post-startup append followed the replacement symlink into capture"
2792 );
2793 assert!(held_root.join("config.json").is_file());
2794 drop(store);
2795 Ok(())
2796 }
2797
2798 #[cfg(all(unix, not(any(target_os = "linux", target_os = "android"))))]
2799 #[test]
2800 fn path_ambient_haematite_refuses_group_or_world_writable_ancestors()
2801 -> Result<(), Box<dyn std::error::Error>> {
2802 use std::os::unix::fs::PermissionsExt as _;
2803
2804 let sandbox = crate::test_support::private_tempdir()?;
2805 std::fs::set_permissions(sandbox.path(), std::fs::Permissions::from_mode(0o700))?;
2806
2807 for mode in [0o770, 0o1777] {
2808 let shared = sandbox.path().join(format!("shared-{mode:o}"));
2809 let data_root = shared.join("data");
2810 std::fs::create_dir(&shared)?;
2811 std::fs::set_permissions(&shared, std::fs::Permissions::from_mode(mode))?;
2812 std::fs::create_dir(&data_root)?;
2813 std::fs::set_permissions(&data_root, std::fs::Permissions::from_mode(0o700))?;
2814 let configured = data_root
2815 .to_str()
2816 .ok_or("temporary data path was not UTF-8")?;
2817
2818 let Err(error) =
2819 super::build_haematite_store(configured, 4, None, test_node_cache_budget()?)
2820 else {
2821 return Err(format!("mode {mode:04o} ancestor was accepted").into());
2822 };
2823 let message = error.to_string();
2824 let crate::ServerError::UnsafeDataRootAncestor {
2825 data_root: resolved_root,
2826 component,
2827 reason,
2828 } = error
2829 else {
2830 return Err(format!("expected typed unsafe-ancestor error, got {message}").into());
2831 };
2832 assert_eq!(resolved_root, std::fs::canonicalize(&data_root)?);
2833 assert_eq!(component, std::fs::canonicalize(&shared)?);
2834 assert!(
2835 reason.contains(&format!("mode {mode:04o}")),
2836 "unexpected reason: {reason}"
2837 );
2838 if mode & 0o1000 != 0 {
2839 assert!(reason.contains("sticky bit is not accepted"));
2840 }
2841 assert!(message.contains("private Aion home"));
2842 assert!(
2843 !data_root.join("config.json").exists(),
2844 "Haematite touched its ambient path before the refusal"
2845 );
2846 }
2847 Ok(())
2848 }
2849
2850 #[cfg(target_os = "macos")]
2851 #[test]
2852 fn path_ambient_haematite_refuses_mutating_allow_acl_ancestor()
2853 -> Result<(), Box<dyn std::error::Error>> {
2854 use std::os::unix::fs::PermissionsExt as _;
2855
2856 let sandbox = crate::test_support::private_tempdir()?;
2857 std::fs::set_permissions(sandbox.path(), std::fs::Permissions::from_mode(0o700))?;
2858 let shared = sandbox.path().join("acl-shared");
2859 let data_root = shared.join("data");
2860 std::fs::create_dir(&shared)?;
2861 std::fs::set_permissions(&shared, std::fs::Permissions::from_mode(0o700))?;
2862 let acl = "everyone allow list,search,add_file,add_subdirectory,delete_child";
2863 let status = std::process::Command::new("chmod")
2864 .arg("+a")
2865 .arg(acl)
2866 .arg(&shared)
2867 .status()?;
2868 assert!(status.success(), "failed to install Darwin regression ACL");
2869 let configured = data_root
2870 .to_str()
2871 .ok_or("temporary data path was not UTF-8")?;
2872
2873 let result = super::build_haematite_store(configured, 4, None, test_node_cache_budget()?);
2874 let cleanup = std::process::Command::new("chmod")
2875 .arg("-RN")
2876 .arg(&shared)
2877 .status()?;
2878 assert!(cleanup.success(), "failed to clean Darwin regression ACL");
2879
2880 let Err(error) = result else {
2881 return Err("mutating non-euid allow ACL ancestor was accepted".into());
2882 };
2883 let message = error.to_string();
2884 let crate::ServerError::UnsafeDataRootAncestor {
2885 component, reason, ..
2886 } = error
2887 else {
2888 return Err(format!("expected typed unsafe-ancestor error, got {message}").into());
2889 };
2890 assert_eq!(component, std::fs::canonicalize(&shared)?);
2891 assert!(
2892 reason.contains("allow"),
2893 "reason did not name the ACE: {reason}"
2894 );
2895 assert!(
2896 reason.contains("everyone"),
2897 "reason did not name the ACE principal: {reason}"
2898 );
2899 assert!(
2900 !data_root.join("config.json").exists(),
2901 "Haematite touched its ambient path before the ACL refusal"
2902 );
2903 Ok(())
2904 }
2905
2906 #[cfg(target_os = "macos")]
2907 #[test]
2908 fn path_ambient_haematite_accepts_the_euid_uuid_allow_ace()
2909 -> Result<(), Box<dyn std::error::Error>> {
2910 use std::os::unix::fs::PermissionsExt as _;
2911
2912 use exacl::{AclEntry, AclOption, Perm};
2913
2914 let sandbox = crate::test_support::private_tempdir()?;
2915 std::fs::set_permissions(sandbox.path(), std::fs::Permissions::from_mode(0o700))?;
2916 let private_parent = sandbox.path().join("euid-uuid-allow");
2917 let data_root = private_parent.join("data");
2918 std::fs::create_dir(&private_parent)?;
2919 std::fs::set_permissions(&private_parent, std::fs::Permissions::from_mode(0o700))?;
2920
2921 let server_uid = rustix::process::geteuid().as_raw();
2922 let ace_qualifier = crate::filesystem::darwin_user_uuid_for_test(server_uid)?;
2923 let entry = AclEntry::allow_user(
2924 &ace_qualifier.to_string(),
2925 Perm::EXECUTE | Perm::WRITE | Perm::APPEND | Perm::DELETE_CHILD,
2926 None,
2927 );
2928 exacl::setfacl(
2929 &[private_parent.as_path()],
2930 &[entry],
2931 AclOption::SYMLINK_ACL,
2932 )?;
2933 let configured = data_root
2934 .to_str()
2935 .ok_or("temporary data path was not UTF-8")?;
2936
2937 let result = super::build_haematite_store(configured, 4, None, test_node_cache_budget()?);
2938 let cleanup = std::process::Command::new("chmod")
2939 .arg("-RN")
2940 .arg(&private_parent)
2941 .status()?;
2942 assert!(cleanup.success(), "failed to clean euid UUID allow ACL");
2943
2944 let (store, responder) = result?;
2945 assert!(responder.is_none());
2946 assert!(data_root.join("config.json").is_file());
2947 drop(store);
2948 Ok(())
2949 }
2950
2951 #[cfg(target_os = "macos")]
2952 #[test]
2953 fn path_ambient_haematite_refuses_a_non_euid_user_uuid_allow_ace()
2954 -> Result<(), Box<dyn std::error::Error>> {
2955 use std::os::unix::fs::PermissionsExt as _;
2956
2957 use exacl::{AclEntry, AclOption, Perm};
2958
2959 let sandbox = crate::test_support::private_tempdir()?;
2960 std::fs::set_permissions(sandbox.path(), std::fs::Permissions::from_mode(0o700))?;
2961 let shared = sandbox.path().join("non-euid-uuid-allow");
2962 let data_root = shared.join("data");
2963 std::fs::create_dir(&shared)?;
2964 std::fs::set_permissions(&shared, std::fs::Permissions::from_mode(0o700))?;
2965
2966 let server_uid = rustix::process::geteuid().as_raw();
2967 let foreign_uid = u32::from(server_uid == 0);
2968 let foreign_qualifier = crate::filesystem::darwin_user_uuid_for_test(foreign_uid)?;
2969 let entry = AclEntry::allow_user(
2970 &foreign_qualifier.to_string(),
2971 Perm::EXECUTE | Perm::WRITE | Perm::APPEND | Perm::DELETE_CHILD,
2972 None,
2973 );
2974 exacl::setfacl(&[shared.as_path()], &[entry], AclOption::SYMLINK_ACL)?;
2975 let configured = data_root
2976 .to_str()
2977 .ok_or("temporary data path was not UTF-8")?;
2978
2979 let result = super::build_haematite_store(configured, 4, None, test_node_cache_budget()?);
2980 let cleanup = std::process::Command::new("chmod")
2981 .arg("-RN")
2982 .arg(&shared)
2983 .status()?;
2984 assert!(cleanup.success(), "failed to clean non-euid UUID allow ACL");
2985
2986 let Err(error) = result else {
2987 return Err("mutating non-euid user UUID allow ACE was accepted".into());
2988 };
2989 let message = error.to_string();
2990 let crate::ServerError::UnsafeDataRootAncestor {
2991 component, reason, ..
2992 } = error
2993 else {
2994 return Err(format!("expected typed unsafe-ancestor error, got {message}").into());
2995 };
2996 assert_eq!(component, std::fs::canonicalize(&shared)?);
2997 assert!(
2998 reason.contains("allow") && reason.contains(&format!("server euid {server_uid}")),
2999 "reason did not name the rejected ACE: {reason}"
3000 );
3001 assert!(
3002 !data_root.join("config.json").exists(),
3003 "Haematite touched its ambient path before the UUID ACL refusal"
3004 );
3005 Ok(())
3006 }
3007
3008 #[cfg(target_os = "macos")]
3009 #[test]
3010 fn path_ambient_haematite_accepts_a_deny_only_acl_ancestor()
3011 -> Result<(), Box<dyn std::error::Error>> {
3012 use std::os::unix::fs::PermissionsExt as _;
3013
3014 let sandbox = crate::test_support::private_tempdir()?;
3015 std::fs::set_permissions(sandbox.path(), std::fs::Permissions::from_mode(0o700))?;
3016 let private_parent = sandbox.path().join("deny-only");
3017 let data_root = private_parent.join("data");
3018 std::fs::create_dir(&private_parent)?;
3019 std::fs::set_permissions(&private_parent, std::fs::Permissions::from_mode(0o700))?;
3020 let status = std::process::Command::new("chmod")
3021 .arg("+a")
3022 .arg("everyone deny delete")
3023 .arg(&private_parent)
3024 .status()?;
3025 assert!(status.success(), "failed to install Darwin deny-only ACL");
3026 let configured = data_root
3027 .to_str()
3028 .ok_or("temporary data path was not UTF-8")?;
3029
3030 let result = super::build_haematite_store(configured, 4, None, test_node_cache_budget()?);
3031 let cleanup = std::process::Command::new("chmod")
3032 .arg("-RN")
3033 .arg(&private_parent)
3034 .status()?;
3035 assert!(cleanup.success(), "failed to clean Darwin deny-only ACL");
3036
3037 let (store, responder) = result?;
3038 assert!(responder.is_none());
3039 assert!(data_root.join("config.json").is_file());
3040 drop(store);
3041 Ok(())
3042 }
3043
3044 #[cfg(target_os = "macos")]
3045 #[test]
3046 fn path_ambient_haematite_accepts_the_stock_home_acl_chain()
3047 -> Result<(), Box<dyn std::error::Error>> {
3048 use std::os::unix::fs::PermissionsExt as _;
3049 use users::os::unix::UserExt as _;
3050
3051 let effective_uid = rustix::process::geteuid().as_raw();
3052 let effective_user = users::get_user_by_uid(effective_uid)
3053 .ok_or_else(|| format!("server euid {effective_uid} has no account record"))?;
3054 let sandbox = tempfile::Builder::new()
3055 .prefix(".aion-acl-home-proof-")
3056 .tempdir_in(effective_user.home_dir())?;
3057 std::fs::set_permissions(sandbox.path(), std::fs::Permissions::from_mode(0o700))?;
3058 let data_root = sandbox.path().join("data");
3059 let configured = data_root
3060 .to_str()
3061 .ok_or("temporary data path was not UTF-8")?;
3062
3063 let (store, responder) =
3064 super::build_haematite_store(configured, 4, None, test_node_cache_budget()?)?;
3065 assert!(responder.is_none());
3066 assert!(data_root.join("config.json").is_file());
3067 drop(store);
3068 Ok(())
3069 }
3070
3071 #[cfg(all(unix, not(any(target_os = "linux", target_os = "android"))))]
3072 #[test]
3073 fn path_ambient_haematite_accepts_an_owner_controlled_chain()
3074 -> Result<(), Box<dyn std::error::Error>> {
3075 use std::os::unix::fs::PermissionsExt as _;
3076
3077 let sandbox = crate::test_support::private_tempdir()?;
3078 std::fs::set_permissions(sandbox.path(), std::fs::Permissions::from_mode(0o700))?;
3079 let private_parent = sandbox.path().join("private");
3080 let data_root = private_parent.join("data");
3081 std::fs::create_dir(&private_parent)?;
3082 std::fs::set_permissions(&private_parent, std::fs::Permissions::from_mode(0o700))?;
3083 let configured = data_root
3084 .to_str()
3085 .ok_or("temporary data path was not UTF-8")?;
3086
3087 let (store, responder) =
3088 super::build_haematite_store(configured, 4, None, test_node_cache_budget()?)?;
3089 assert!(responder.is_none());
3090 assert!(data_root.join("config.json").is_file());
3091 for shard in 0..4 {
3092 assert!(data_root.join(format!("shard-{shard}")).is_dir());
3093 }
3094 drop(store);
3095 Ok(())
3096 }
3097
3098 #[tokio::test]
3099 async fn connect_store_memory_backend_exposes_no_outbox_store()
3100 -> Result<(), Box<dyn std::error::Error>> {
3101 use crate::config::{StoreBackend, StoreConfig};
3102
3103 let connected = super::connect_store(StoreConfig {
3106 backend: StoreBackend::Memory,
3107 owned_shards: Vec::new(),
3108 data_dir: None,
3109 shard_count: 1,
3110 cluster: None,
3111 node_cache_budget: None,
3112 ..StoreConfig::default()
3113 })
3114 .await?;
3115 assert!(
3116 connected.outbox_store.is_none(),
3117 "the in-memory backend exposes no outbox store"
3118 );
3119 Ok(())
3120 }
3121
3122 #[tokio::test]
3128 async fn haematite_boot_refuses_without_a_node_cache_budget()
3129 -> Result<(), Box<dyn std::error::Error>> {
3130 use crate::ServerError;
3131 use crate::config::{StoreBackend, StoreConfig};
3132
3133 let sandbox = crate::test_support::private_tempdir()?;
3134 let data_dir = sandbox.path().join("data");
3135 let error = super::connect_haematite_store(StoreConfig {
3136 backend: StoreBackend::Haematite,
3137 data_dir: Some(
3138 data_dir
3139 .to_str()
3140 .ok_or("temporary data path was not UTF-8")?
3141 .to_owned(),
3142 ),
3143 shard_count: 4,
3144 ..StoreConfig::default()
3145 })
3146 .await
3147 .err()
3148 .ok_or("the haematite boot path must refuse a store config with no node_cache_budget")?;
3149 let ServerError::Config { message } = error else {
3150 return Err(format!("expected a config refusal, got {error:?}").into());
3151 };
3152 assert!(
3153 message.contains("store.node_cache_budget"),
3154 "the refusal must name the missing key, got: {message}"
3155 );
3156 assert!(
3157 message.contains("AION_STORE_NODE_CACHE_BUDGET"),
3158 "the refusal must name the environment override, got: {message}"
3159 );
3160 Ok(())
3161 }
3162
3163 #[tokio::test]
3168 async fn configured_node_cache_budget_reaches_the_created_database()
3169 -> Result<(), Box<dyn std::error::Error>> {
3170 use crate::config::{StoreBackend, StoreConfig};
3171
3172 const ONE_GIB: usize = 1 << 30;
3173
3174 let sandbox = crate::test_support::private_tempdir()?;
3175 let data_dir = sandbox.path().join("data");
3176 let connected = super::connect_haematite_store(StoreConfig {
3177 backend: StoreBackend::Haematite,
3178 data_dir: Some(
3179 data_dir
3180 .to_str()
3181 .ok_or("temporary data path was not UTF-8")?
3182 .to_owned(),
3183 ),
3184 shard_count: 4,
3185 node_cache_budget: Some(haematite::NodeCacheBudget::bytes(ONE_GIB)?),
3186 ..StoreConfig::default()
3187 })
3188 .await?;
3189 drop(connected);
3190
3191 let recorded: serde_json::Value =
3192 serde_json::from_slice(&std::fs::read(data_dir.join("config.json"))?)?;
3193 assert_eq!(
3194 recorded.get("node_cache_budget"),
3195 Some(&serde_json::json!({ "bytes": ONE_GIB })),
3196 "the operator's budget must be the one haematite created the database with"
3197 );
3198 Ok(())
3199 }
3200
3201 #[tokio::test]
3202 async fn state_build_fails_without_event_broadcast_capacity()
3203 -> Result<(), Box<dyn std::error::Error>> {
3204 let mut runtime = runtime_config();
3205 runtime.websocket.event_broadcast_capacity = None;
3206
3207 let error = ServerState::build_with_store(InMemoryStore::default(), runtime)
3208 .await
3209 .err()
3210 .ok_or("state build must fail when event streaming is unsized")?;
3211
3212 assert!(error.is_config(), "expected a config error, got {error}");
3213 assert!(
3214 error
3215 .to_string()
3216 .contains("websocket.event_broadcast_capacity"),
3217 "error must name the missing key: {error}"
3218 );
3219 Ok(())
3220 }
3221
3222 #[tokio::test]
3223 async fn state_build_fails_without_query_timeout() -> Result<(), Box<dyn std::error::Error>> {
3224 let mut runtime = runtime_config();
3225 runtime.query_timeout = None;
3226
3227 let error = ServerState::build_with_store(InMemoryStore::default(), runtime)
3228 .await
3229 .err()
3230 .ok_or("state build must fail when the query reply deadline is unset")?;
3231
3232 assert!(error.is_config(), "expected a config error, got {error}");
3233 assert!(
3234 error.to_string().contains("runtime.query_timeout_ms"),
3235 "error must name the missing key: {error}"
3236 );
3237 assert!(
3238 error.to_string().contains("AION_RUNTIME_QUERY_TIMEOUT_MS"),
3239 "error must name the environment override: {error}"
3240 );
3241 Ok(())
3242 }
3243
3244 #[tokio::test]
3245 async fn state_build_fails_with_zero_query_timeout() -> Result<(), Box<dyn std::error::Error>> {
3246 let mut runtime = runtime_config();
3247 runtime.query_timeout = Some(Duration::ZERO);
3248
3249 let error = ServerState::build_with_store(InMemoryStore::default(), runtime)
3250 .await
3251 .err()
3252 .ok_or("state build must fail when the query reply deadline is zero")?;
3253
3254 assert!(error.is_config(), "expected a config error, got {error}");
3255 assert!(
3256 error.to_string().contains("runtime.query_timeout_ms"),
3257 "error must name the zero-valued key: {error}"
3258 );
3259 Ok(())
3260 }
3261
3262 #[tokio::test(flavor = "multi_thread")]
3285 async fn a_completed_check_through_the_built_dispatcher_lands_in_the_returned_slot()
3286 -> Result<(), Box<dyn std::error::Error>> {
3287 use std::collections::{BTreeMap, VecDeque};
3288 use std::sync::{Arc, Mutex};
3289
3290 use aion::ActivityDispatch;
3291 use aion_core::{ActivityId, RunId, WorkflowId};
3292 use aion_package::ActionBodyContract;
3293
3294 use crate::update_check::document::{FETCH_ACTION, FETCH_COMMAND, UPDATE_CHECK_QUEUE};
3295 use crate::worker::{DeclaredBodies, DeclaredBodyLookup, DispatchingRun};
3296
3297 struct SequencedBodies {
3300 replies: Mutex<VecDeque<DeclaredBodyLookup>>,
3301 }
3302
3303 impl DeclaredBodies for SequencedBodies {
3304 fn body_for(
3305 &self,
3306 _task_queue: &str,
3307 _action: &str,
3308 _run: DispatchingRun<'_>,
3309 ) -> DeclaredBodyLookup {
3310 let mut replies = match self.replies.lock() {
3311 Ok(replies) => replies,
3312 Err(poisoned) => poisoned.into_inner(),
3313 };
3314 replies.pop_front().unwrap_or(DeclaredBodyLookup::None)
3315 }
3316 }
3317
3318 let runtime = runtime_config();
3319 let cluster_publisher = crate::cluster_publisher::ClusterEventPublisher::new(
3320 ServerState::FALLBACK_CLUSTER_BROADCAST_CAPACITY,
3321 );
3322 let namespace_store: Arc<dyn aion_store::NamespaceStore> =
3323 Arc::new(InMemoryStore::default());
3324 let worker_deployment_store: Arc<dyn aion_store::WorkerDeploymentStore> =
3325 Arc::new(InMemoryStore::default());
3326 let seams = super::build_worker_seams(
3327 &runtime,
3328 &cluster_publisher,
3329 &namespace_store,
3330 &worker_deployment_store,
3331 None,
3332 );
3333
3334 let fixture = concat!(
3337 env!("CARGO_MANIFEST_DIR"),
3338 "/src/update_check/fixtures/aion-cli-index.jsonl"
3339 );
3340 seams.declared_bodies.install(Arc::new(SequencedBodies {
3341 replies: Mutex::new(VecDeque::from([
3342 DeclaredBodyLookup::Declared(ActionBodyContract::Run {
3343 command: FETCH_COMMAND.to_owned(),
3344 }),
3345 DeclaredBodyLookup::Declared(ActionBodyContract::Run {
3346 command: format!("cat {fixture}"),
3347 }),
3348 ])),
3349 }));
3350
3351 let transcript = crate::activity_publisher::ActivityEventPublisher::new(
3352 Arc::new(aion_store::InMemoryObservabilityStore::default()),
3353 ServerState::FALLBACK_CLUSTER_BROADCAST_CAPACITY,
3354 crate::activity_publisher::TranscriptBatchPolicy::UNBATCHED,
3355 );
3356 let (dispatcher, _mock_registry, _attempt_owners, _workspace_root, update_status) =
3357 super::build_decorated_dispatcher(&runtime, &seams, transcript);
3358
3359 assert_eq!(
3360 update_status.last(),
3361 None,
3362 "the returned slot must start honestly empty"
3363 );
3364
3365 let dispatch = ActivityDispatch {
3366 namespace: "default".to_owned(),
3367 task_queue: UPDATE_CHECK_QUEUE.to_owned(),
3368 node: None,
3369 workflow_id: WorkflowId::new_v4(),
3370 run_id: RunId::new_v4(),
3371 activity_id: ActivityId::from_sequence_position(1),
3372 name: FETCH_ACTION.to_owned(),
3373 input: "{}".to_owned(),
3374 config: "{}".to_owned(),
3375 attempt: 1,
3376 labels: BTreeMap::new(),
3377 advisory: false,
3378 };
3379 let handle = tokio::task::spawn_blocking(move || dispatcher.dispatch(dispatch));
3380 let encoded = handle
3381 .await?
3382 .map_err(|error| format!("the check dispatch failed: {error}"))?;
3383 let outcome: serde_json::Value = serde_json::from_str(&encoded)?;
3384 assert_eq!(outcome["exit_code"], 0, "the local stand-in command ran");
3385
3386 let recorded = update_status.last().ok_or(
3387 "the completed check must land in the RETURNED slot — the one the boot path \
3388 stores and /update-status serves; an empty slot here is the disconnected-\
3389 producer mis-wire the r1 review proved unmeasured",
3390 )?;
3391 assert_eq!(recorded.latest_known, "0.13.7");
3392 Ok(())
3393 }
3394}