1mod boot_required;
8
9pub(crate) use boot_required::{
10 BOOT_REQUIRED_FIELD_DEFAULTS, NOT_UPGRADE_HEALABLE, RequiredFieldDefault, boot_required_probe,
11};
12
13use std::{path::PathBuf, sync::Arc};
14
15use aion::{
16 ActivityDispatcher, EngineBuilder, RuntimeHandle, SignalRouter, signal::ConcreteSignalRouter,
17};
18use aion_store::{EventStore, NamespaceStore, OutboxStore, WorkerDeploymentStore};
19
20use crate::dev_ui::{ActivityMockRegistry, DevMockingDispatcher};
21
22#[cfg(feature = "auth")]
23use crate::auth::JwksCache;
24use crate::{
25 config::{RuntimeConfig, ServerConfig, StoreBackend, StoreConfig},
26 error::ServerError,
27 namespace::{NamespaceGuard, NamespaceMinter, resolver::NamespaceResolver},
28 observability::{
29 Metrics, health::HealthState, instrumented_store::InstrumentedEventStore,
30 metrics::MetricsError,
31 },
32 shutdown::DrainState,
33 worker::{
34 ConnectedWorkerRegistry, HeartbeatTracker, PendingActivities, WorkerActivityDispatcher,
35 supervisor::WorkerSupervisor,
36 },
37};
38
39fn assistant_mcp_endpoint(
61 runtime: &RuntimeConfig,
62) -> Option<crate::assistant::sessions::launch::AssistantEndpoints> {
63 let configured = runtime.listen.http;
64 if configured.port() == 0 {
65 tracing::warn!(
66 "`server.listen_address` names port 0, so this server cannot state an address for an \
67 assistant agent to dial before it has bound; assistant sessions will be opened with \
68 no MCP server of ours — including the `assistant_context` tool, so an agent will not \
69 be able to read what is on the operator's screen"
70 );
71 return None;
72 }
73 let host = if configured.ip().is_unspecified() {
74 std::net::IpAddr::V4(std::net::Ipv4Addr::LOCALHOST)
75 } else {
76 configured.ip()
77 };
78 let authority = std::net::SocketAddr::new(host, configured.port());
79 Some(crate::assistant::sessions::launch::AssistantEndpoints {
80 base: format!("http://{authority}"),
81 aion_mcp_enabled: runtime.mcp.enabled,
82 })
83}
84
85fn build_assistant_sessions(
89 store: Arc<dyn aion_store::assistant::AssistantSessionStore>,
90 runtime: &RuntimeConfig,
91) -> crate::assistant::sessions::AssistantSessions {
92 crate::assistant::sessions::AssistantSessions::new(
93 store,
94 runtime.assistant.clone(),
95 assistant_mcp_endpoint(runtime),
96 )
97}
98
99fn new_supervisor(
110 store: &Arc<dyn WorkerDeploymentStore>,
111 publisher: &crate::cluster_publisher::ClusterEventPublisher,
112) -> Arc<WorkerSupervisor> {
113 Arc::new(WorkerSupervisor::new(Arc::clone(store), publisher.clone()))
114}
115
116#[derive(Clone)]
118pub struct ServerState {
119 inner: Arc<ServerStateInner>,
120}
121
122struct ServerStateInner {
123 namespace_guard: NamespaceGuard,
124 runtime: RuntimeConfig,
125 worker_registry: ConnectedWorkerRegistry,
126 pending_activities: PendingActivities,
127 heartbeat_tracker: HeartbeatTracker,
128 grpc_liveness_waiters: crate::worker::GrpcLivenessWaiters,
135 drain_state: DrainState,
136 metrics: Option<Metrics>,
137 health: Option<HealthState>,
138 activity_mock_registry: Option<ActivityMockRegistry>,
141 outbox_store: Option<Arc<dyn OutboxStore>>,
144 namespace_store: Arc<dyn NamespaceStore>,
153 worker_deployment_store: Arc<dyn WorkerDeploymentStore>,
155 worker_supervisor: Arc<WorkerSupervisor>,
160 outbox_wake: Arc<tokio::sync::Notify>,
167 cluster_publisher: crate::cluster_publisher::ClusterEventPublisher,
171 transcript_publisher: crate::activity_publisher::ActivityEventPublisher,
180 attempt_owners: crate::worker::AttemptOwnerIndex,
187 queue_service_state: crate::worker::QueueServiceState,
193 queue_declarations: crate::worker::QueueDeclarationSource,
197 declared_attempts: crate::worker::DeclaredCommandAttempts,
204 declared_bodies: crate::worker::DeclaredBodySource,
209 update_status: crate::update_check::UpdateStatusState,
215 workspace_root: crate::worker::WorkspaceRoot,
223 cluster_self_node: Option<String>,
228 cluster_responder: Option<aion_store_haematite::ClusterResponder>,
233 cluster_store: Option<Arc<aion_store_haematite::HaematiteStore>>,
236 watched_peers: Vec<crate::cluster::WatchedPeer>,
239 shard_directory: Option<Arc<crate::routing::StaticShardDirectory>>,
244 request_forwarder: Option<Arc<dyn crate::routing::RequestForwarder>>,
248 assistant_sessions: crate::assistant::sessions::AssistantSessions,
256 #[cfg(feature = "auth")]
257 jwks_cache: Option<JwksCache>,
258}
259
260impl ServerState {
261 const FALLBACK_CLUSTER_BROADCAST_CAPACITY: std::num::NonZeroUsize =
270 match std::num::NonZeroUsize::new(64) {
271 Some(value) => value,
272 None => std::num::NonZeroUsize::MIN,
273 };
274
275 pub async fn build(
282 config: ServerConfig,
283 stage: &crate::control::StageReporter,
284 ) -> Result<Self, ServerError> {
285 let (store_config, runtime) = config.into_parts();
286 let connected = connect_store(store_config, stage).await?;
287 Self::build_with_connected_store(connected, runtime, stage).await
288 }
289
290 pub async fn build_with_store<S>(store: S, runtime: RuntimeConfig) -> Result<Self, ServerError>
296 where
297 S: EventStore
298 + NamespaceStore
299 + WorkerDeploymentStore
300 + aion_store::workloop::WorkloopStore
301 + aion_store::visibility::VisibilityStore
302 + aion_store::AssistantSessionStore,
303 {
304 let leaf = Arc::new(store);
309 let namespace_store: Arc<dyn NamespaceStore> = leaf.clone();
310 let worker_deployment_store: Arc<dyn WorkerDeploymentStore> = leaf.clone();
311 let workloop_store: Arc<dyn aion_store::workloop::WorkloopStore> = leaf.clone();
312 let visibility_store: Arc<dyn aion_store::visibility::VisibilityStore> = leaf.clone();
313 let assistant_store: Arc<dyn aion_store::AssistantSessionStore> = leaf.clone();
314 Self::build_with_connected_store(
318 ConnectedStore::local(
319 leaf,
320 None,
321 namespace_store,
322 worker_deployment_store,
323 workloop_store,
324 visibility_store,
325 assistant_store,
326 ),
327 runtime,
328 &crate::control::StageReporter::detached(),
329 )
330 .await
331 }
332
333 async fn build_with_connected_store(
334 connected: ConnectedStore,
335 runtime: RuntimeConfig,
336 stage: &crate::control::StageReporter,
337 ) -> Result<Self, ServerError> {
338 let cluster_self_node = connected.cluster_self_node();
339 let outbox_store = connected.outbox_store;
340 let bootstrap_coordinator = connected.bootstrap_coordinator;
341 let cluster_responder = connected.cluster_responder;
342 let cluster_store = connected.cluster_store;
343 let watched_peers = connected.watched_peers;
344 let RoutingState {
348 shard_directory,
349 request_forwarder,
350 mint_routing,
351 } = build_routing_state(
352 cluster_store.as_ref(),
353 connected.directory_peers,
354 connected.self_node_id,
355 );
356 let (event_broadcast_capacity, query_timeout, workloop_sweep_interval) =
357 required_engine_seams(&runtime)?;
358 let (cluster_publisher, transcript_publisher) =
359 build_real_time_publishers(&runtime, connected.observability_store)?;
360 let (metrics, outbox_wake, instrumented_store) =
361 build_instrumented_store(&runtime, connected.event_store)?;
362 let exported_metrics = runtime.metrics.enabled.then_some(metrics.clone());
363 let seams = build_worker_seams(
364 &runtime,
365 &cluster_publisher,
366 &metrics,
367 &connected.namespace_store,
368 &connected.worker_deployment_store,
369 mint_routing,
370 );
371 let (
372 activity_dispatcher,
373 activity_mock_registry,
374 attempt_owners,
375 workspace_root,
376 update_status,
377 ) = build_decorated_dispatcher(&runtime, &seams, transcript_publisher.clone());
378
379 let engine = boot_engine(EngineAssembly {
380 seams: &seams,
381 instrumented_store: &instrumented_store,
382 event_broadcast_capacity,
383 query_timeout,
384 workloop_store: Arc::clone(&connected.workloop_store),
385 visibility_store: Arc::clone(&connected.visibility_store),
386 workloop_sweep_interval,
387 activity_dispatcher,
388 active_registry: Arc::new(aion::Registry::default()),
389 bootstrap_coordinator,
390 runtime: &runtime,
391 stage,
392 })
393 .await?;
394 let resolver = NamespaceResolver::from_config(runtime.namespace.clone(), engine);
395 let worker_supervisor =
396 new_supervisor(&connected.worker_deployment_store, &cluster_publisher);
397 #[cfg(feature = "auth")]
398 let jwks_cache = build_jwks_cache(&runtime).await?;
399 let assistant_sessions =
400 build_assistant_sessions(Arc::clone(&connected.assistant_store), &runtime);
401 Ok(Self {
402 inner: Arc::new(ServerStateInner {
403 namespace_guard: NamespaceGuard::new(resolver),
404 assistant_sessions,
405 runtime,
406 metrics: exported_metrics,
407 worker_registry: seams.worker_registry,
408 pending_activities: seams.pending_activities,
409 heartbeat_tracker: seams.heartbeat_tracker,
410 grpc_liveness_waiters: crate::worker::GrpcLivenessWaiters::new(),
411 drain_state: seams.drain_state,
412 health: Some(HealthState::new(instrumented_store, true)),
413 activity_mock_registry,
414 outbox_store,
415 namespace_store: connected.namespace_store,
416 worker_supervisor,
417 worker_deployment_store: connected.worker_deployment_store,
418 outbox_wake,
419 cluster_publisher,
420 transcript_publisher,
421 attempt_owners,
422 queue_service_state: seams.queue_service_state,
423 queue_declarations: seams.queue_declarations,
424 declared_attempts: seams.declared_attempts,
425 declared_bodies: seams.declared_bodies.clone(),
426 update_status,
427 workspace_root,
428 cluster_self_node,
429 cluster_responder,
430 cluster_store,
431 watched_peers,
432 shard_directory,
433 request_forwarder,
434 #[cfg(feature = "auth")]
435 jwks_cache,
436 }),
437 })
438 }
439
440 #[must_use]
442 pub fn from_parts(namespace_resolver: NamespaceResolver, runtime: RuntimeConfig) -> Self {
443 Self::from_parts_with_namespace_store(
448 namespace_resolver,
449 runtime,
450 Arc::new(aion_store::InMemoryStore::default()),
451 )
452 }
453
454 #[must_use]
461 pub fn from_parts_with_namespace_store<S>(
462 namespace_resolver: NamespaceResolver,
463 runtime: RuntimeConfig,
464 store: Arc<S>,
465 ) -> Self
466 where
467 S: NamespaceStore + WorkerDeploymentStore,
468 {
469 let namespace_store: Arc<dyn NamespaceStore> = store.clone();
470 let worker_deployment_store: Arc<dyn WorkerDeploymentStore> = store;
471 Self::from_parts_with_control_stores(
472 namespace_resolver,
473 runtime,
474 namespace_store,
475 worker_deployment_store,
476 )
477 }
478
479 #[must_use]
484 pub fn from_parts_with_control_stores(
485 namespace_resolver: NamespaceResolver,
486 runtime: RuntimeConfig,
487 namespace_store: Arc<dyn NamespaceStore>,
488 worker_deployment_store: Arc<dyn WorkerDeploymentStore>,
489 ) -> Self {
490 let heartbeat_tracker = HeartbeatTracker::new(runtime.worker.heartbeat_window);
491 let pending_activities = PendingActivities::new(runtime.worker.heartbeat_window);
497 let bounds = transcript_bounds(&runtime);
501 let batch = required_transcript_batch_policy(&runtime)
505 .unwrap_or(crate::activity_publisher::TranscriptBatchPolicy::UNBATCHED);
506 let cluster_publisher = crate::cluster_publisher::ClusterEventPublisher::new(
507 Self::FALLBACK_CLUSTER_BROADCAST_CAPACITY,
508 );
509 let drain_state = DrainState::default();
510 let assistant_sessions =
511 build_assistant_sessions(Arc::new(aion_store::InMemoryStore::default()), &runtime);
512 Self {
513 inner: Arc::new(ServerStateInner {
514 namespace_guard: NamespaceGuard::new(namespace_resolver),
515 assistant_sessions,
516 runtime,
517 worker_registry: ConnectedWorkerRegistry::default()
518 .with_worker_deployment_store(worker_deployment_store.clone())
519 .with_cluster_publisher(cluster_publisher.clone()),
520 pending_activities,
521 heartbeat_tracker,
522 grpc_liveness_waiters: crate::worker::GrpcLivenessWaiters::new(),
523 drain_state: drain_state.clone(),
524 metrics: None,
525 health: None,
526 activity_mock_registry: None,
527 outbox_store: None,
528 namespace_store,
529 worker_supervisor: new_supervisor(&worker_deployment_store, &cluster_publisher),
530 worker_deployment_store,
531 outbox_wake: Arc::new(tokio::sync::Notify::new()),
532 cluster_publisher,
533 transcript_publisher: build_transcript_publisher(
537 None,
538 Self::FALLBACK_CLUSTER_BROADCAST_CAPACITY,
539 bounds,
540 batch,
541 ),
542 attempt_owners: crate::worker::AttemptOwnerIndex::new(),
543 queue_service_state: crate::worker::QueueServiceState::default(),
544 queue_declarations: crate::worker::QueueDeclarationSource::default(),
545 declared_attempts: crate::worker::DeclaredCommandAttempts::new(drain_state.clone()),
549 declared_bodies: crate::worker::DeclaredBodySource::default(),
550 update_status: crate::update_check::UpdateStatusState::default(),
554 workspace_root: crate::worker::WorkspaceRoot::resolve(),
558 cluster_self_node: None,
559 cluster_responder: None,
560 cluster_store: None,
561 watched_peers: Vec::new(),
562 shard_directory: None,
563 request_forwarder: None,
564 #[cfg(feature = "auth")]
565 jwks_cache: None,
566 }),
567 }
568 }
569
570 #[cfg(feature = "auth")]
579 #[must_use]
580 pub fn from_parts_with_namespace_store_and_jwks<S>(
581 namespace_resolver: NamespaceResolver,
582 runtime: RuntimeConfig,
583 store: Arc<S>,
584 jwks_cache: JwksCache,
585 ) -> Self
586 where
587 S: NamespaceStore + WorkerDeploymentStore,
588 {
589 let namespace_store: Arc<dyn NamespaceStore> = store.clone();
590 let worker_deployment_store: Arc<dyn WorkerDeploymentStore> = store;
591 let heartbeat_tracker = HeartbeatTracker::new(runtime.worker.heartbeat_window);
592 let pending_activities = PendingActivities::new(runtime.worker.heartbeat_window);
598 let bounds = transcript_bounds(&runtime);
602 let batch = required_transcript_batch_policy(&runtime)
606 .unwrap_or(crate::activity_publisher::TranscriptBatchPolicy::UNBATCHED);
607 let cluster_publisher = crate::cluster_publisher::ClusterEventPublisher::new(
608 Self::FALLBACK_CLUSTER_BROADCAST_CAPACITY,
609 );
610 let drain_state = DrainState::default();
611 let assistant_sessions =
612 build_assistant_sessions(Arc::new(aion_store::InMemoryStore::default()), &runtime);
613 Self {
614 inner: Arc::new(ServerStateInner {
615 namespace_guard: NamespaceGuard::new(namespace_resolver),
616 assistant_sessions,
617 runtime,
618 worker_registry: ConnectedWorkerRegistry::default()
619 .with_worker_deployment_store(worker_deployment_store.clone())
620 .with_cluster_publisher(cluster_publisher.clone()),
621 pending_activities,
622 heartbeat_tracker,
623 grpc_liveness_waiters: crate::worker::GrpcLivenessWaiters::new(),
624 drain_state: drain_state.clone(),
625 metrics: None,
626 health: None,
627 activity_mock_registry: None,
628 outbox_store: None,
629 namespace_store,
630 worker_supervisor: new_supervisor(&worker_deployment_store, &cluster_publisher),
631 worker_deployment_store,
632 outbox_wake: Arc::new(tokio::sync::Notify::new()),
633 cluster_publisher,
634 transcript_publisher: build_transcript_publisher(
638 None,
639 Self::FALLBACK_CLUSTER_BROADCAST_CAPACITY,
640 bounds,
641 batch,
642 ),
643 attempt_owners: crate::worker::AttemptOwnerIndex::new(),
644 queue_service_state: crate::worker::QueueServiceState::default(),
645 queue_declarations: crate::worker::QueueDeclarationSource::default(),
646 declared_attempts: crate::worker::DeclaredCommandAttempts::new(drain_state.clone()),
650 declared_bodies: crate::worker::DeclaredBodySource::default(),
651 update_status: crate::update_check::UpdateStatusState::default(),
655 workspace_root: crate::worker::WorkspaceRoot::resolve(),
659 cluster_self_node: None,
660 cluster_responder: None,
661 cluster_store: None,
662 watched_peers: Vec::new(),
663 shard_directory: None,
664 request_forwarder: None,
665 jwks_cache: Some(jwks_cache),
666 }),
667 }
668 }
669
670 #[cfg(feature = "auth")]
676 #[must_use]
677 pub fn from_parts_with_jwks(
678 namespace_resolver: NamespaceResolver,
679 runtime: RuntimeConfig,
680 jwks_cache: JwksCache,
681 ) -> Self {
682 let fallback_deployment_store: Arc<dyn WorkerDeploymentStore> =
686 Arc::new(aion_store::InMemoryStore::default());
687 let heartbeat_tracker = HeartbeatTracker::new(runtime.worker.heartbeat_window);
688 let pending_activities = PendingActivities::new(runtime.worker.heartbeat_window);
694 let bounds = transcript_bounds(&runtime);
698 let batch = required_transcript_batch_policy(&runtime)
702 .unwrap_or(crate::activity_publisher::TranscriptBatchPolicy::UNBATCHED);
703 let cluster_publisher = crate::cluster_publisher::ClusterEventPublisher::new(
704 Self::FALLBACK_CLUSTER_BROADCAST_CAPACITY,
705 );
706 let drain_state = DrainState::default();
707 let assistant_sessions =
708 build_assistant_sessions(Arc::new(aion_store::InMemoryStore::default()), &runtime);
709 Self {
710 inner: Arc::new(ServerStateInner {
711 namespace_guard: NamespaceGuard::new(namespace_resolver),
712 assistant_sessions,
713 runtime,
714 worker_registry: ConnectedWorkerRegistry::default(),
715 pending_activities,
716 heartbeat_tracker,
717 grpc_liveness_waiters: crate::worker::GrpcLivenessWaiters::new(),
718 drain_state: drain_state.clone(),
719 metrics: None,
720 health: None,
721 activity_mock_registry: None,
722 outbox_store: None,
723 namespace_store: Arc::new(aion_store::InMemoryStore::default()),
728 worker_supervisor: new_supervisor(&fallback_deployment_store, &cluster_publisher),
729 worker_deployment_store: fallback_deployment_store,
730 outbox_wake: Arc::new(tokio::sync::Notify::new()),
731 cluster_publisher,
732 transcript_publisher: build_transcript_publisher(
736 None,
737 Self::FALLBACK_CLUSTER_BROADCAST_CAPACITY,
738 bounds,
739 batch,
740 ),
741 attempt_owners: crate::worker::AttemptOwnerIndex::new(),
742 queue_service_state: crate::worker::QueueServiceState::default(),
743 queue_declarations: crate::worker::QueueDeclarationSource::default(),
744 declared_attempts: crate::worker::DeclaredCommandAttempts::new(drain_state.clone()),
748 declared_bodies: crate::worker::DeclaredBodySource::default(),
749 update_status: crate::update_check::UpdateStatusState::default(),
753 workspace_root: crate::worker::WorkspaceRoot::resolve(),
757 cluster_self_node: None,
758 cluster_responder: None,
759 cluster_store: None,
760 watched_peers: Vec::new(),
761 shard_directory: None,
762 request_forwarder: None,
763 jwks_cache: Some(jwks_cache),
764 }),
765 }
766 }
767
768 #[must_use]
770 pub fn from_parts_with_registry(
771 namespace_resolver: NamespaceResolver,
772 runtime: RuntimeConfig,
773 worker_registry: ConnectedWorkerRegistry,
774 ) -> Self {
775 let fallback_deployment_store: Arc<dyn WorkerDeploymentStore> =
779 Arc::new(aion_store::InMemoryStore::default());
780 let heartbeat_tracker = HeartbeatTracker::new(runtime.worker.heartbeat_window);
781 let pending_activities = PendingActivities::new(runtime.worker.heartbeat_window);
787 let bounds = transcript_bounds(&runtime);
791 let batch = required_transcript_batch_policy(&runtime)
795 .unwrap_or(crate::activity_publisher::TranscriptBatchPolicy::UNBATCHED);
796 let cluster_publisher = crate::cluster_publisher::ClusterEventPublisher::new(
797 Self::FALLBACK_CLUSTER_BROADCAST_CAPACITY,
798 );
799 let drain_state = DrainState::default();
800 let assistant_sessions =
801 build_assistant_sessions(Arc::new(aion_store::InMemoryStore::default()), &runtime);
802 Self {
803 inner: Arc::new(ServerStateInner {
804 namespace_guard: NamespaceGuard::new(namespace_resolver),
805 assistant_sessions,
806 runtime,
807 worker_registry,
808 pending_activities,
809 heartbeat_tracker,
810 grpc_liveness_waiters: crate::worker::GrpcLivenessWaiters::new(),
811 drain_state: drain_state.clone(),
812 metrics: None,
813 health: None,
814 activity_mock_registry: None,
815 outbox_store: None,
816 namespace_store: Arc::new(aion_store::InMemoryStore::default()),
821 worker_supervisor: new_supervisor(&fallback_deployment_store, &cluster_publisher),
822 worker_deployment_store: fallback_deployment_store,
823 outbox_wake: Arc::new(tokio::sync::Notify::new()),
824 cluster_publisher,
825 transcript_publisher: build_transcript_publisher(
829 None,
830 Self::FALLBACK_CLUSTER_BROADCAST_CAPACITY,
831 bounds,
832 batch,
833 ),
834 attempt_owners: crate::worker::AttemptOwnerIndex::new(),
835 queue_service_state: crate::worker::QueueServiceState::default(),
836 queue_declarations: crate::worker::QueueDeclarationSource::default(),
837 declared_attempts: crate::worker::DeclaredCommandAttempts::new(drain_state.clone()),
841 declared_bodies: crate::worker::DeclaredBodySource::default(),
842 update_status: crate::update_check::UpdateStatusState::default(),
846 workspace_root: crate::worker::WorkspaceRoot::resolve(),
850 cluster_self_node: None,
851 cluster_responder: None,
852 cluster_store: None,
853 watched_peers: Vec::new(),
854 shard_directory: None,
855 request_forwarder: None,
856 #[cfg(feature = "auth")]
857 jwks_cache: None,
858 }),
859 }
860 }
861
862 #[must_use]
864 pub fn namespace_guard(&self) -> &NamespaceGuard {
865 &self.inner.namespace_guard
866 }
867
868 #[must_use]
870 pub fn deploy_guard(&self) -> crate::deploy::DeployGuard {
871 crate::deploy::DeployGuard::new(self.inner.namespace_guard.resolver().clone())
872 }
873
874 #[must_use]
876 pub fn runtime_config(&self) -> &RuntimeConfig {
877 &self.inner.runtime
878 }
879
880 #[must_use]
888 pub(crate) fn update_status(&self) -> &crate::update_check::UpdateStatusState {
889 &self.inner.update_status
890 }
891
892 #[must_use]
902 pub fn workspace_root(&self) -> &crate::worker::WorkspaceRoot {
903 &self.inner.workspace_root
904 }
905
906 #[must_use]
914 pub fn assistant_sessions(&self) -> &crate::assistant::sessions::AssistantSessions {
915 &self.inner.assistant_sessions
916 }
917
918 #[must_use]
920 pub fn worker_registry(&self) -> &ConnectedWorkerRegistry {
921 &self.inner.worker_registry
922 }
923
924 pub fn cancel_in_flight_activities(
943 &self,
944 workflow_id: &aion_core::WorkflowId,
945 ) -> Result<crate::worker::InFlightCancellation, ServerError> {
946 crate::worker::cancel_in_flight_activities(
947 self.heartbeat_tracker(),
948 self.worker_registry(),
949 self.declared_attempts(),
950 workflow_id,
951 )
952 }
953
954 #[must_use]
959 pub fn declared_attempts(&self) -> &crate::worker::DeclaredCommandAttempts {
960 &self.inner.declared_attempts
961 }
962
963 #[must_use]
968 pub fn declared_bodies(&self) -> &crate::worker::DeclaredBodySource {
969 &self.inner.declared_bodies
970 }
971
972 #[must_use]
976 pub fn cluster_publisher(&self) -> &crate::cluster_publisher::ClusterEventPublisher {
977 &self.inner.cluster_publisher
978 }
979
980 #[must_use]
985 pub fn transcript_publisher(&self) -> &crate::activity_publisher::ActivityEventPublisher {
986 &self.inner.transcript_publisher
987 }
988
989 #[must_use]
993 pub fn attempt_owners(&self) -> &crate::worker::AttemptOwnerIndex {
994 &self.inner.attempt_owners
995 }
996
997 #[must_use]
1000 pub fn queue_service_state(&self) -> &crate::worker::QueueServiceState {
1001 &self.inner.queue_service_state
1002 }
1003
1004 #[must_use]
1008 pub fn queue_declarations(&self) -> &crate::worker::QueueDeclarationSource {
1009 &self.inner.queue_declarations
1010 }
1011
1012 pub fn unserved_queues(&self) -> Result<Vec<crate::worker::UnservedQueue>, ServerError> {
1023 self.inner.queue_service_state.unserved()
1024 }
1025
1026 pub fn unrecoverable_runs(
1044 &self,
1045 ) -> Result<Vec<(aion_core::WorkflowId, aion::registry::UnrecoverableRun)>, ServerError> {
1046 self.engine()?
1047 .registry()
1048 .unrecoverable()
1049 .list()
1050 .map_err(ServerError::from)
1051 }
1052
1053 #[must_use]
1065 pub fn intervention_router(&self) -> crate::worker::InterventionRouter {
1066 let transport: std::sync::Arc<dyn crate::worker::InterventionTransport> = {
1067 #[cfg(feature = "liminal-transport")]
1068 {
1069 std::sync::Arc::new(crate::worker::LiminalInterventionTransport)
1070 }
1071 #[cfg(not(feature = "liminal-transport"))]
1072 {
1073 std::sync::Arc::new(NullInterventionTransport)
1074 }
1075 };
1076 crate::worker::InterventionRouter::new(
1077 self.inner.worker_registry.clone(),
1078 self.inner.attempt_owners.clone(),
1079 transport,
1080 )
1081 .with_transcript_publisher(self.inner.transcript_publisher.clone())
1084 }
1085
1086 #[must_use]
1090 pub fn cluster_self_node(&self) -> Option<&str> {
1091 self.inner.cluster_self_node.as_deref()
1092 }
1093
1094 pub fn engine(&self) -> Result<Arc<aion::Engine>, ServerError> {
1107 self.inner
1108 .namespace_guard
1109 .resolver()
1110 .engine()
1111 .map(Arc::clone)
1112 }
1113
1114 #[must_use]
1116 pub fn pending_activities(&self) -> &PendingActivities {
1117 &self.inner.pending_activities
1118 }
1119
1120 #[must_use]
1124 pub fn read_provenance(&self) -> aion_core::ReadProvenance {
1125 aion_core::ReadProvenance::new(
1126 self.inner
1127 .pending_activities
1128 .lease_recorder()
1129 .ledger()
1130 .failures(),
1131 )
1132 }
1133
1134 #[must_use]
1136 pub fn heartbeat_tracker(&self) -> &HeartbeatTracker {
1137 &self.inner.heartbeat_tracker
1138 }
1139
1140 #[must_use]
1146 pub fn grpc_liveness_waiters(&self) -> &crate::worker::GrpcLivenessWaiters {
1147 &self.inner.grpc_liveness_waiters
1148 }
1149
1150 #[must_use]
1152 pub fn drain_state(&self) -> &DrainState {
1153 &self.inner.drain_state
1154 }
1155
1156 #[must_use]
1158 pub fn metrics(&self) -> Option<&Metrics> {
1159 self.inner.metrics.as_ref()
1160 }
1161
1162 #[must_use]
1164 pub fn health(&self) -> Option<&HealthState> {
1165 self.inner.health.as_ref()
1166 }
1167
1168 #[must_use]
1173 pub fn activity_mock_registry(&self) -> Option<&ActivityMockRegistry> {
1174 self.inner.activity_mock_registry.as_ref()
1175 }
1176
1177 #[must_use]
1183 pub fn outbox_store(&self) -> Option<Arc<dyn OutboxStore>> {
1184 self.inner.outbox_store.clone()
1185 }
1186
1187 #[must_use]
1196 pub fn namespace_store(&self) -> &Arc<dyn NamespaceStore> {
1197 &self.inner.namespace_store
1198 }
1199
1200 #[must_use]
1202 pub fn worker_deployment_store(&self) -> &Arc<dyn WorkerDeploymentStore> {
1203 &self.inner.worker_deployment_store
1204 }
1205
1206 #[must_use]
1213 pub fn worker_supervisor(&self) -> &Arc<WorkerSupervisor> {
1214 &self.inner.worker_supervisor
1215 }
1216
1217 #[must_use]
1226 pub fn namespace_minter(&self) -> NamespaceMinter {
1227 let minter = NamespaceMinter::new(
1228 Arc::clone(&self.inner.namespace_store),
1229 self.inner.runtime.auto_create,
1230 )
1231 .with_cluster_publisher(self.inner.cluster_publisher.clone());
1236 match self.namespace_routing() {
1241 Some(routing) => minter.with_routing(routing),
1242 None => minter,
1243 }
1244 }
1245
1246 #[must_use]
1256 pub fn namespace_routing(&self) -> Option<crate::namespace::NamespaceRouting> {
1257 build_namespace_routing(
1258 self.cluster_store(),
1259 self.shard_directory(),
1260 self.request_forwarder(),
1261 )
1262 }
1263
1264 #[must_use]
1270 pub fn outbox_wake(&self) -> Arc<tokio::sync::Notify> {
1271 Arc::clone(&self.inner.outbox_wake)
1272 }
1273
1274 #[must_use]
1280 pub fn is_clustered(&self) -> bool {
1281 self.inner.cluster_responder.is_some()
1282 }
1283
1284 #[must_use]
1289 pub fn cluster_store(&self) -> Option<&Arc<aion_store_haematite::HaematiteStore>> {
1290 self.inner.cluster_store.as_ref()
1291 }
1292
1293 #[must_use]
1297 pub fn shard_directory(&self) -> Option<&Arc<crate::routing::StaticShardDirectory>> {
1298 self.inner.shard_directory.as_ref()
1299 }
1300
1301 #[must_use]
1304 pub fn request_forwarder(&self) -> Option<&Arc<dyn crate::routing::RequestForwarder>> {
1305 self.inner.request_forwarder.as_ref()
1306 }
1307
1308 #[must_use]
1325 pub fn spawn_heartbeat_sweeper(
1326 &self,
1327 shutdown: tokio::sync::watch::Receiver<bool>,
1328 ) -> tokio::task::JoinHandle<()> {
1329 let sweeper = crate::worker::HeartbeatSweeper::new(
1330 self.inner.heartbeat_tracker.clone(),
1331 self.inner.worker_registry.clone(),
1332 self.inner.pending_activities.clone(),
1333 self.inner.drain_state.clone(),
1334 self.inner.runtime.worker.heartbeat_window,
1335 )
1336 .with_queue_state(self.inner.queue_service_state.clone());
1340 tokio::spawn(sweeper.run(shutdown))
1341 }
1342
1343 pub fn spawn_startup_catchup(
1367 &self,
1368 mut shutdown: tokio::sync::watch::Receiver<bool>,
1369 ) -> Result<tokio::task::JoinHandle<()>, ServerError> {
1370 let engine = self.engine()?;
1371 Ok(tokio::spawn(async move {
1372 tracing::info!(
1373 "startup catch-up running behind open doors: owed timer fires, \
1374 schedule catch-up"
1375 );
1376 let catchup = engine.run_startup_catchup();
1377 tokio::pin!(catchup);
1378 tokio::select! {
1379 result = &mut catchup => match result {
1380 Ok(()) => {
1381 tracing::info!("startup catch-up complete");
1382 }
1383 Err(error) => {
1384 tracing::error!(
1385 %error,
1386 "startup catch-up failed; owed timer fires and schedule \
1387 catch-up remain undelivered — restart the server to re-run \
1388 the sweep from durable state"
1389 );
1390 }
1391 },
1392 _ = shutdown.changed() => {
1393 tracing::info!(
1394 "startup catch-up interrupted by shutdown; the next boot's \
1395 sweep resumes from durable state"
1396 );
1397 }
1398 }
1399 }))
1400 }
1401
1402 #[cfg(feature = "liminal-transport")]
1413 #[must_use]
1414 pub fn spawn_liminal_liveness_probe(
1415 &self,
1416 notifier: std::sync::Arc<crate::worker::LiminalConnectionNotifier>,
1417 shutdown: tokio::sync::watch::Receiver<bool>,
1418 ) -> tokio::task::JoinHandle<()> {
1419 let probe = crate::worker::LivenessProbe::across_transports(
1420 Some(notifier),
1421 Some(self.inner.grpc_liveness_waiters.clone()),
1426 self.inner.heartbeat_tracker.clone(),
1427 self.inner.worker_registry.clone(),
1428 self.inner.runtime.worker.heartbeat_window,
1429 );
1430 tokio::spawn(probe.run(shutdown))
1431 }
1432
1433 pub fn spawn_cluster_supervisor(
1449 &self,
1450 config: crate::cluster::SupervisorConfig,
1451 shutdown: tokio::sync::watch::Receiver<bool>,
1452 ) -> Result<bool, ServerError> {
1453 let Some(cluster_store) = self.inner.cluster_store.clone() else {
1454 return Ok(false);
1455 };
1456 if self.inner.watched_peers.is_empty() {
1457 return Ok(false);
1458 }
1459 let engine = Arc::clone(self.inner.namespace_guard.resolver().engine()?);
1460 let publisher = Arc::new(self.inner.cluster_publisher.clone());
1464 let self_node = self.inner.cluster_self_node.clone().unwrap_or_default();
1465 let adopter = Arc::new(crate::cluster::OutboxSettlingAdopter::new(
1471 engine,
1472 self.inner.outbox_store.clone(),
1473 ));
1474 let supervisor = crate::cluster::ClusterSupervisor::new(
1475 cluster_store,
1476 adopter,
1477 self.inner.watched_peers.clone(),
1478 config,
1479 )
1480 .with_publisher(publisher, self_node);
1481 if !supervisor.watches_any() {
1482 return Ok(false);
1483 }
1484 tokio::spawn(supervisor.run(shutdown));
1485 Ok(true)
1486 }
1487
1488 #[cfg(feature = "auth")]
1490 #[must_use]
1491 pub fn jwks_cache(&self) -> Option<&JwksCache> {
1492 self.inner.jwks_cache.as_ref()
1493 }
1494
1495 pub fn shutdown(&self) -> Result<(), ServerError> {
1502 self.inner.namespace_guard.resolver().shutdown_engine()
1503 }
1504}
1505
1506#[cfg(feature = "auth")]
1507async fn build_jwks_cache(runtime: &RuntimeConfig) -> Result<Option<JwksCache>, ServerError> {
1508 if !runtime.auth.enabled {
1509 return Ok(None);
1510 }
1511 let Some(url) = runtime.auth.jwks_url.clone() else {
1512 return Err(ServerError::Config {
1513 message: "auth.jwks_url must not be empty when auth.enabled is true".to_owned(),
1514 });
1515 };
1516 let interval = std::time::Duration::from_secs(runtime.auth.jwks_refresh_seconds);
1517 let cache = JwksCache::new(url, interval)
1518 .await
1519 .map_err(|error| ServerError::Config {
1520 message: format!("auth jwks initial fetch failed: {error}"),
1521 })?;
1522 Ok(Some(cache))
1523}
1524
1525fn metrics_config_error(error: &MetricsError) -> ServerError {
1526 ServerError::Config {
1527 message: error.to_string(),
1528 }
1529}
1530
1531struct EngineAssembly<'a> {
1533 seams: &'a WorkerSeams,
1535 instrumented_store: &'a Arc<InstrumentedEventStore>,
1537 event_broadcast_capacity: std::num::NonZeroUsize,
1539 query_timeout: std::time::Duration,
1541 workloop_store: Arc<dyn aion_store::workloop::WorkloopStore>,
1545 visibility_store: Arc<dyn aion_store::visibility::VisibilityStore>,
1549 workloop_sweep_interval: std::time::Duration,
1552 activity_dispatcher: Arc<dyn ActivityDispatcher>,
1554 active_registry: Arc<aion::Registry>,
1556 bootstrap_coordinator: bool,
1558 runtime: &'a RuntimeConfig,
1560 stage: &'a crate::control::StageReporter,
1564}
1565
1566fn server_search_attribute_schema() -> Result<aion_core::SearchAttributeSchema, ServerError> {
1575 let mut schema = aion_core::SearchAttributeSchema::new();
1576 for (name, label) in [
1577 (crate::namespace::NAMESPACE_ATTRIBUTE, "namespace"),
1578 (crate::namespace::TASK_QUEUE_ATTRIBUTE, "task_queue"),
1579 (crate::namespace::DISPLAY_NAME_ATTRIBUTE, "display_name"),
1580 ] {
1581 schema
1582 .register(name, aion_core::SearchAttributeType::String)
1583 .map_err(|error| ServerError::Config {
1584 message: format!("failed to register {label} search attribute: {error}"),
1585 })?;
1586 }
1587 Ok(schema)
1588}
1589
1590async fn boot_engine(assembly: EngineAssembly<'_>) -> Result<Arc<aion::Engine>, ServerError> {
1600 let search_attribute_schema = server_search_attribute_schema()?;
1601 let runtime = assembly.runtime;
1602 let builder = EngineBuilder::new()
1603 .store_arc(assembly.instrumented_store.clone())
1604 .event_streaming(assembly.event_broadcast_capacity)
1605 .visibility_store_arc(assembly.visibility_store)
1606 .search_attribute_schema(search_attribute_schema)
1607 .scheduler_threads(runtime.scheduler_threads)
1608 .stop_drain_timeout(
1612 runtime
1613 .stop_drain_timeout
1614 .ok_or_else(|| ServerError::Config {
1615 message: String::from(crate::config::STOP_DRAIN_TIMEOUT_REQUIRED),
1616 })?,
1617 )
1618 .outbox_enabled(runtime.outbox.enabled)
1619 .activity_dispatcher(assembly.activity_dispatcher)
1620 .active_registry(assembly.active_registry)
1621 .production_recovery_seam()
1622 .defer_startup_recovery()
1629 .signal_router_factory(|runtime: Arc<RuntimeHandle>, handoff| {
1630 Arc::new(ConcreteSignalRouter::new(runtime, handoff)) as Arc<dyn SignalRouter>
1631 })
1632 .query_timeout(assembly.query_timeout)
1633 .with_workloop_service(
1641 Arc::clone(&assembly.workloop_store),
1642 assembly.workloop_sweep_interval,
1643 )
1644 .bootstrap_schedule_coordinator(assembly.bootstrap_coordinator)
1650 .load_workflow_sources(runtime.workflow_packages.iter().map(PathBuf::as_path));
1651 let builder = if runtime.owned_shards.is_empty() {
1657 builder
1658 } else {
1659 builder.owned_shards(runtime.owned_shards.iter().copied())
1660 };
1661 let builder = match runtime.jit_threshold {
1667 Some(threshold) => builder.scheduler_jit_threshold(threshold),
1668 None => builder,
1669 };
1670 let engine = Arc::new(builder.build().await.map_err(ServerError::from)?);
1671 install_engine_backed_seams(assembly.seams, &engine, runtime.outbox.enabled);
1672 assembly.stage.report(
1685 crate::control::stage::STAGE_ENGINE_RECOVERY,
1686 "recovering resident workflows from durable state".to_owned(),
1687 );
1688 engine
1689 .recover_workflows_on_startup()
1690 .await
1691 .map_err(ServerError::from)?;
1692 assembly.stage.report(
1693 crate::control::stage::STAGE_ENGINE_RECOVERY,
1694 "resident workflows recovered".to_owned(),
1695 );
1696 Ok(engine)
1697}
1698
1699fn required_engine_seams(
1704 runtime: &RuntimeConfig,
1705) -> Result<
1706 (
1707 std::num::NonZeroUsize,
1708 std::time::Duration,
1709 std::time::Duration,
1710 ),
1711 ServerError,
1712> {
1713 let event_broadcast_capacity = runtime
1714 .websocket
1715 .event_broadcast_capacity
1716 .and_then(std::num::NonZeroUsize::new)
1717 .ok_or_else(|| ServerError::Config {
1718 message: crate::config::EVENT_BROADCAST_CAPACITY_REQUIRED.to_owned(),
1719 })?;
1720 let query_timeout = runtime
1721 .query_timeout
1722 .filter(|timeout| !timeout.is_zero())
1723 .ok_or_else(|| ServerError::Config {
1724 message: crate::config::QUERY_TIMEOUT_REQUIRED.to_owned(),
1725 })?;
1726 let workloop_sweep_interval = runtime
1731 .workloop_sweep_interval
1732 .filter(|interval| !interval.is_zero())
1733 .ok_or_else(|| ServerError::Config {
1734 message: crate::config::WORKLOOP_SWEEP_INTERVAL_REQUIRED.to_owned(),
1735 })?;
1736 Ok((
1737 event_broadcast_capacity,
1738 query_timeout,
1739 workloop_sweep_interval,
1740 ))
1741}
1742
1743fn install_outbox_delivery(
1750 pending_activities: &PendingActivities,
1751 engine: &Arc<aion::Engine>,
1752 outbox_enabled: bool,
1753) {
1754 if outbox_enabled {
1755 let callback = Arc::new(crate::worker::ServerOutboxDeliveryCallback::new(
1756 Arc::clone(engine),
1757 ));
1758 pending_activities.set_outbox_delivery(callback);
1759 }
1760}
1761
1762struct WorkerSeams {
1764 worker_registry: ConnectedWorkerRegistry,
1765 pending_activities: PendingActivities,
1766 heartbeat_tracker: HeartbeatTracker,
1767 drain_state: DrainState,
1768 queue_declarations: crate::worker::QueueDeclarationSource,
1769 queue_service_state: crate::worker::QueueServiceState,
1770 declared_bodies: crate::worker::DeclaredBodySource,
1771 declared_attempts: crate::worker::DeclaredCommandAttempts,
1777 metrics: Metrics,
1780 cluster_publisher: crate::cluster_publisher::ClusterEventPublisher,
1784}
1785
1786fn build_worker_seams(
1791 runtime: &RuntimeConfig,
1792 cluster_publisher: &crate::cluster_publisher::ClusterEventPublisher,
1793 metrics: &Metrics,
1794 namespace_store: &Arc<dyn NamespaceStore>,
1795 worker_deployment_store: &Arc<dyn WorkerDeploymentStore>,
1796 mint_routing: Option<crate::namespace::NamespaceRouting>,
1797) -> WorkerSeams {
1798 let worker_registry = ConnectedWorkerRegistry::default()
1799 .with_cluster_publisher(cluster_publisher.clone())
1800 .with_namespace_minting(namespace_store.clone(), runtime.auto_create)
1801 .with_worker_deployment_store(worker_deployment_store.clone());
1802 let worker_registry = match mint_routing {
1808 Some(routing) => worker_registry.with_namespace_routing(routing),
1809 None => worker_registry,
1810 };
1811 let drain_state = DrainState::default();
1812 WorkerSeams {
1813 worker_registry,
1814 pending_activities: PendingActivities::new(runtime.worker.heartbeat_window),
1819 heartbeat_tracker: HeartbeatTracker::new(runtime.worker.heartbeat_window),
1820 declared_attempts: crate::worker::DeclaredCommandAttempts::new(drain_state.clone()),
1821 drain_state,
1822 queue_declarations: crate::worker::QueueDeclarationSource::default(),
1823 queue_service_state: crate::worker::QueueServiceState::default(),
1824 declared_bodies: crate::worker::DeclaredBodySource::default(),
1825 metrics: metrics.clone(),
1826 cluster_publisher: cluster_publisher.clone(),
1827 }
1828}
1829
1830fn install_queue_declarations(
1837 queue_declarations: &crate::worker::QueueDeclarationSource,
1838 engine: &Arc<aion::Engine>,
1839) {
1840 queue_declarations.install(Arc::new(crate::worker::EngineQueueDeclarations::new(
1841 Arc::clone(engine),
1842 )));
1843}
1844
1845fn install_declared_bodies(
1860 declared_bodies: &crate::worker::DeclaredBodySource,
1861 engine: &Arc<aion::Engine>,
1862) {
1863 declared_bodies.install(Arc::new(crate::worker::EngineDeclaredBodies::new(
1864 Arc::clone(engine),
1865 )));
1866}
1867
1868fn install_engine_backed_seams(seams: &WorkerSeams, engine: &Arc<aion::Engine>, outbox: bool) {
1874 install_outbox_delivery(&seams.pending_activities, engine, outbox);
1875 seams.pending_activities.set_lease_recorder(
1879 Arc::new(crate::worker::EngineLeaseRecorder::new(Arc::clone(engine))),
1880 Some(seams.metrics.clone()),
1881 );
1882 install_queue_declarations(&seams.queue_declarations, engine);
1883 install_declared_bodies(&seams.declared_bodies, engine);
1884}
1885
1886fn build_bridge_dispatcher(
1899 runtime: &RuntimeConfig,
1900 seams: &WorkerSeams,
1901) -> (WorkerActivityDispatcher, crate::worker::AttemptOwnerIndex) {
1902 let attempt_owners = crate::worker::AttemptOwnerIndex::new();
1903 let dispatcher = WorkerActivityDispatcher::new(
1904 seams.worker_registry.clone(),
1905 runtime.default_namespace.clone(),
1906 seams.heartbeat_tracker.clone(),
1907 )
1908 .with_pending(seams.pending_activities.clone())
1909 .with_drain_state(seams.drain_state.clone())
1910 .with_tokio_handle(tokio::runtime::Handle::current())
1911 .with_attempt_owners(attempt_owners.clone())
1912 .with_queue_service(runtime.worker.queue_service.clone())
1913 .with_queue_declarations(seams.queue_declarations.clone())
1914 .with_queue_state(seams.queue_service_state.clone())
1915 .with_cluster_publisher(seams.cluster_publisher.clone());
1916 (dispatcher, attempt_owners)
1917}
1918
1919fn build_instrumented_store(
1932 runtime: &RuntimeConfig,
1933 event_store: Arc<dyn EventStore>,
1934) -> Result<
1935 (
1936 Metrics,
1937 Arc<tokio::sync::Notify>,
1938 Arc<InstrumentedEventStore>,
1939 ),
1940 ServerError,
1941> {
1942 let metrics = Metrics::new().map_err(|error| metrics_config_error(&error))?;
1943 let outbox_wake = Arc::new(tokio::sync::Notify::new());
1944 let instrumented_store = Arc::new(
1945 InstrumentedEventStore::new(
1946 event_store,
1947 metrics.clone(),
1948 runtime.default_namespace.clone(),
1949 )
1950 .with_outbox_wake(Arc::clone(&outbox_wake)),
1951 );
1952 Ok((metrics, outbox_wake, instrumented_store))
1953}
1954
1955fn build_decorated_dispatcher(
1960 runtime: &RuntimeConfig,
1961 seams: &WorkerSeams,
1962 transcript: crate::activity_publisher::ActivityEventPublisher,
1963) -> (
1964 Arc<dyn ActivityDispatcher>,
1965 Option<ActivityMockRegistry>,
1966 crate::worker::AttemptOwnerIndex,
1967 crate::worker::WorkspaceRoot,
1968 crate::update_check::UpdateStatusState,
1969) {
1970 let workspace_root = crate::worker::WorkspaceRoot::resolve();
1974 let update_status = crate::update_check::UpdateStatusState::default();
1977 let (dispatcher, attempt_owners) = build_bridge_dispatcher(runtime, seams);
1978 let (activity_dispatcher, activity_mock_registry) = decorate_activity_dispatcher(
1979 dispatcher,
1980 seams.declared_bodies.clone(),
1981 seams.declared_attempts.clone(),
1982 workspace_root.clone(),
1983 transcript,
1984 update_status.clone(),
1985 runtime.dev.enabled,
1986 );
1987 (
1988 activity_dispatcher,
1989 activity_mock_registry,
1990 attempt_owners,
1991 workspace_root,
1992 update_status,
1993 )
1994}
1995
1996fn decorate_activity_dispatcher(
1997 dispatcher: WorkerActivityDispatcher,
1998 declared_bodies: crate::worker::DeclaredBodySource,
1999 declared_attempts: crate::worker::DeclaredCommandAttempts,
2000 workspace_root: crate::worker::WorkspaceRoot,
2001 transcript: crate::activity_publisher::ActivityEventPublisher,
2002 update_status: crate::update_check::UpdateStatusState,
2003 dev_enabled: bool,
2004) -> (Arc<dyn ActivityDispatcher>, Option<ActivityMockRegistry>) {
2005 let declared = crate::worker::DeclaredCommandDispatcher::new(
2013 Arc::new(dispatcher),
2014 declared_bodies.clone(),
2015 declared_attempts,
2016 tokio::runtime::Handle::current(),
2017 workspace_root,
2018 transcript,
2019 );
2020 let observed = crate::update_check::UpdateCheckObserver::new(
2021 Arc::new(declared),
2022 declared_bodies,
2023 update_status,
2024 );
2025 if dev_enabled {
2026 let registry = ActivityMockRegistry::new();
2027 let decorated = DevMockingDispatcher::new(Arc::new(observed), registry.clone());
2028 (Arc::new(decorated), Some(registry))
2029 } else {
2030 (Arc::new(observed), None)
2031 }
2032}
2033
2034fn required_cluster_broadcast_capacity(
2039 runtime: &RuntimeConfig,
2040) -> Result<std::num::NonZeroUsize, ServerError> {
2041 runtime
2042 .websocket
2043 .cluster_broadcast_capacity
2044 .and_then(std::num::NonZeroUsize::new)
2045 .ok_or_else(|| ServerError::Config {
2046 message: crate::config::CLUSTER_BROADCAST_CAPACITY_REQUIRED.to_owned(),
2047 })
2048}
2049
2050fn build_real_time_publishers(
2063 runtime: &RuntimeConfig,
2064 observability_store: Option<Arc<dyn aion_store::ObservabilityStore>>,
2065) -> Result<
2066 (
2067 crate::cluster_publisher::ClusterEventPublisher,
2068 crate::activity_publisher::ActivityEventPublisher,
2069 ),
2070 ServerError,
2071> {
2072 let capacity = required_cluster_broadcast_capacity(runtime)?;
2073 let batch = required_transcript_batch_policy(runtime)?;
2074 Ok((
2075 crate::cluster_publisher::ClusterEventPublisher::new(capacity),
2076 build_transcript_publisher(
2077 observability_store,
2078 capacity,
2079 transcript_bounds(runtime),
2080 batch,
2081 ),
2082 ))
2083}
2084
2085pub(crate) fn required_node_cache_budget(
2104 config: &StoreConfig,
2105) -> Result<haematite::NodeCacheBudget, ServerError> {
2106 config.node_cache_budget.ok_or_else(|| ServerError::Config {
2107 message: crate::config::STORE_NODE_CACHE_BUDGET_REQUIRED.to_owned(),
2108 })
2109}
2110
2111fn warn_retired_lock_acquisition_keys(config: &StoreConfig) {
2118 if config.lock_acquisition_patience_ms.is_some()
2119 || config.lock_acquisition_retry_cadence_ms.is_some()
2120 {
2121 tracing::warn!(
2122 "store.lock_acquisition_patience_ms and store.lock_acquisition_retry_cadence_ms \
2123 are retired and ignored: the server waits indefinitely for the data-directory \
2124 writer lock (a kernel-released flock whose holder is either live mid-handover \
2125 or already gone) and reports while it waits. Remove the keys"
2126 );
2127 }
2128}
2129
2130pub(crate) fn required_transcript_batch_policy(
2150 runtime: &RuntimeConfig,
2151) -> Result<crate::activity_publisher::TranscriptBatchPolicy, ServerError> {
2152 let max_batch_events = runtime
2153 .observability
2154 .max_batch_events
2155 .and_then(std::num::NonZeroUsize::new)
2156 .ok_or_else(|| ServerError::Config {
2157 message: crate::config::OBSERVABILITY_MAX_BATCH_EVENTS_REQUIRED.to_owned(),
2158 })?;
2159 let max_batch_hold_ms =
2162 runtime
2163 .observability
2164 .max_batch_hold_ms
2165 .ok_or_else(|| ServerError::Config {
2166 message: crate::config::OBSERVABILITY_MAX_BATCH_HOLD_MS_REQUIRED.to_owned(),
2167 })?;
2168 Ok(crate::activity_publisher::TranscriptBatchPolicy {
2169 max_batch_events,
2170 max_hold: std::time::Duration::from_millis(max_batch_hold_ms),
2171 })
2172}
2173
2174fn transcript_bounds(runtime: &RuntimeConfig) -> crate::activity_bounds::TranscriptBounds {
2176 crate::activity_bounds::TranscriptBounds {
2177 max_event_bytes: runtime.observability.max_event_bytes,
2178 max_stream_events: runtime.observability.max_stream_events,
2179 }
2180}
2181
2182fn build_transcript_publisher(
2193 observability_store: Option<Arc<dyn aion_store::ObservabilityStore>>,
2194 capacity: std::num::NonZeroUsize,
2195 bounds: crate::activity_bounds::TranscriptBounds,
2196 batch: crate::activity_publisher::TranscriptBatchPolicy,
2197) -> crate::activity_publisher::ActivityEventPublisher {
2198 let store = observability_store
2199 .unwrap_or_else(|| Arc::new(aion_store::InMemoryObservabilityStore::default()));
2200 crate::activity_publisher::ActivityEventPublisher::new(store, capacity, batch)
2201 .with_bounds(bounds)
2202}
2203
2204struct RoutingState {
2206 shard_directory: Option<Arc<crate::routing::StaticShardDirectory>>,
2207 request_forwarder: Option<Arc<dyn crate::routing::RequestForwarder>>,
2208 mint_routing: Option<crate::namespace::NamespaceRouting>,
2213}
2214
2215fn build_routing_state(
2219 cluster_store: Option<&Arc<aion_store_haematite::HaematiteStore>>,
2220 directory_peers: Vec<crate::routing::DirectoryPeer>,
2221 self_node_id: Option<String>,
2222) -> RoutingState {
2223 let Some(store) = cluster_store else {
2224 return RoutingState {
2225 shard_directory: None,
2226 request_forwarder: None,
2227 mint_routing: None,
2228 };
2229 };
2230 let shard_directory = Arc::new(crate::routing::StaticShardDirectory::new(
2231 Arc::clone(store),
2232 directory_peers,
2233 self_node_id,
2234 ));
2235 let request_forwarder: Arc<dyn crate::routing::RequestForwarder> =
2236 Arc::new(crate::routing::GrpcRequestForwarder::new());
2237 let mint_routing = build_namespace_routing(
2238 Some(store),
2239 Some(&shard_directory),
2240 Some(&request_forwarder),
2241 );
2242 RoutingState {
2243 shard_directory: Some(shard_directory),
2244 request_forwarder: Some(request_forwarder),
2245 mint_routing,
2246 }
2247}
2248
2249fn build_namespace_routing(
2259 cluster_store: Option<&Arc<aion_store_haematite::HaematiteStore>>,
2260 shard_directory: Option<&Arc<crate::routing::StaticShardDirectory>>,
2261 request_forwarder: Option<&Arc<dyn crate::routing::RequestForwarder>>,
2262) -> Option<crate::namespace::NamespaceRouting> {
2263 use crate::namespace::{
2264 GrpcMintForwarder, MintForwarder, MintShardOwners, NamespaceRouting, NamespaceShardResolver,
2265 };
2266 let store = Arc::clone(cluster_store?);
2267 let directory = Arc::clone(shard_directory?);
2268 let shards: Arc<dyn NamespaceShardResolver> = store;
2269 let owners: Arc<dyn MintShardOwners> = directory;
2270 let forwarder: Arc<dyn MintForwarder> =
2271 Arc::new(GrpcMintForwarder::new(Arc::clone(request_forwarder?)));
2272 Some(NamespaceRouting::new(shards, owners, forwarder))
2273}
2274
2275struct ConnectedStore {
2285 event_store: Arc<dyn EventStore>,
2286 outbox_store: Option<Arc<dyn OutboxStore>>,
2287 namespace_store: Arc<dyn NamespaceStore>,
2294 worker_deployment_store: Arc<dyn WorkerDeploymentStore>,
2296 workloop_store: Arc<dyn aion_store::workloop::WorkloopStore>,
2306 visibility_store: Arc<dyn aion_store::visibility::VisibilityStore>,
2311 assistant_store: Arc<dyn aion_store::AssistantSessionStore>,
2318 observability_store: Option<Arc<dyn aion_store::ObservabilityStore>>,
2326 bootstrap_coordinator: bool,
2327 cluster_responder: Option<aion_store_haematite::ClusterResponder>,
2328 cluster_store: Option<Arc<aion_store_haematite::HaematiteStore>>,
2332 watched_peers: Vec<crate::cluster::WatchedPeer>,
2335 directory_peers: Vec<crate::routing::DirectoryPeer>,
2339 self_node_id: Option<String>,
2343}
2344
2345impl ConnectedStore {
2346 fn local(
2353 event_store: Arc<dyn EventStore>,
2354 outbox_store: Option<Arc<dyn OutboxStore>>,
2355 namespace_store: Arc<dyn NamespaceStore>,
2356 worker_deployment_store: Arc<dyn WorkerDeploymentStore>,
2357 workloop_store: Arc<dyn aion_store::workloop::WorkloopStore>,
2358 visibility_store: Arc<dyn aion_store::visibility::VisibilityStore>,
2359 assistant_store: Arc<dyn aion_store::AssistantSessionStore>,
2360 ) -> Self {
2361 Self {
2362 event_store,
2363 outbox_store,
2364 namespace_store,
2365 worker_deployment_store,
2366 workloop_store,
2367 visibility_store,
2368 assistant_store,
2369 observability_store: None,
2375 bootstrap_coordinator: true,
2376 cluster_responder: None,
2377 cluster_store: None,
2378 watched_peers: Vec::new(),
2379 directory_peers: Vec::new(),
2380 self_node_id: None,
2381 }
2382 }
2383
2384 fn cluster_self_node(&self) -> Option<String> {
2387 self.self_node_id.clone()
2388 }
2389}
2390
2391async fn connect_store(
2394 config: StoreConfig,
2395 stage: &crate::control::StageReporter,
2396) -> Result<ConnectedStore, ServerError> {
2397 match config.backend {
2398 StoreBackend::Memory => {
2399 let leaf = Arc::new(aion_store::InMemoryStore::default());
2402 let namespace_store: Arc<dyn NamespaceStore> = leaf.clone();
2403 let worker_deployment_store: Arc<dyn WorkerDeploymentStore> = leaf.clone();
2404 let workloop_store: Arc<dyn aion_store::workloop::WorkloopStore> = leaf.clone();
2405 let visibility_store: Arc<dyn aion_store::visibility::VisibilityStore> = leaf.clone();
2406 let assistant_store: Arc<dyn aion_store::AssistantSessionStore> = leaf.clone();
2407 Ok(ConnectedStore::local(
2408 leaf,
2409 None,
2410 namespace_store,
2411 worker_deployment_store,
2412 workloop_store,
2413 visibility_store,
2414 assistant_store,
2415 ))
2416 }
2417 StoreBackend::Haematite => connect_haematite_store(config, stage).await,
2418 }
2419}
2420
2421async fn connect_haematite_store(
2441 config: StoreConfig,
2442 stage: &crate::control::StageReporter,
2443) -> Result<ConnectedStore, ServerError> {
2444 let node_cache_budget = required_node_cache_budget(&config)?;
2450 warn_retired_lock_acquisition_keys(&config);
2451 let Some(data_dir) = config.data_dir else {
2452 return Err(ServerError::Config {
2453 message: "store.data_dir must not be empty when store.backend is haematite".to_owned(),
2454 });
2455 };
2456 let shard_count = config.shard_count;
2457 let owned_shards = config.owned_shards.clone();
2458 let cluster = config.cluster.clone();
2459 let watched_peers: Vec<crate::cluster::WatchedPeer> = cluster
2464 .as_ref()
2465 .map(|cluster| {
2466 cluster
2467 .peers
2468 .iter()
2469 .map(|peer| crate::cluster::WatchedPeer {
2470 name: peer.name.clone(),
2471 owned_shards: peer.owned_shards.clone(),
2472 })
2473 .collect()
2474 })
2475 .unwrap_or_default();
2476 let directory_peers: Vec<crate::routing::DirectoryPeer> = cluster
2479 .as_ref()
2480 .map(|cluster| {
2481 cluster
2482 .peers
2483 .iter()
2484 .map(|peer| crate::routing::DirectoryPeer {
2485 name: peer.name.clone(),
2486 owned_shards: peer.owned_shards.clone(),
2487 grpc_addr: peer.grpc_address,
2488 })
2489 .collect()
2490 })
2491 .unwrap_or_default();
2492 let self_node_id: Option<String> = cluster.as_ref().map(|cluster| cluster.node_id.clone());
2495 stage.report(
2499 crate::control::stage::STAGE_STORE_OPEN,
2500 format!("opening the haematite store at {data_dir}"),
2501 );
2502 let build_stage = stage.clone();
2507 let (store, responder) = tokio::task::spawn_blocking(move || {
2508 build_haematite_store(
2509 &data_dir,
2510 shard_count,
2511 cluster,
2512 node_cache_budget,
2513 &build_stage,
2514 )
2515 })
2516 .await
2517 .map_err(|error| ServerError::Config {
2518 message: format!("haematite store initialization task failed: {error}"),
2519 })??;
2520
2521 let bootstrap_coordinator = if owned_shards.is_empty() {
2525 true
2526 } else {
2527 store.set_owned_shards(owned_shards.iter().copied());
2528 store.owns_workflow_shard(&aion::schedule_coordinator_workflow_id())
2529 };
2530
2531 let leaf = Arc::new(store);
2532 let event_store: Arc<dyn EventStore> = leaf.clone();
2533 let outbox_store: Arc<dyn OutboxStore> = leaf.clone();
2534 let namespace_store: Arc<dyn NamespaceStore> = leaf.clone();
2538 let worker_deployment_store: Arc<dyn WorkerDeploymentStore> = leaf.clone();
2539 let workloop_store: Arc<dyn aion_store::workloop::WorkloopStore> = leaf.clone();
2540 let visibility_store: Arc<dyn aion_store::visibility::VisibilityStore> = leaf.clone();
2543 let assistant_store: Arc<dyn aion_store::AssistantSessionStore> = leaf.clone();
2544 let observability_store: Arc<dyn aion_store::ObservabilityStore> = leaf.clone();
2549 let cluster_store = responder.as_ref().map(|_| leaf);
2553 let (watched_peers, directory_peers, self_node_id) = if cluster_store.is_some() {
2554 (watched_peers, directory_peers, self_node_id)
2555 } else {
2556 (Vec::new(), Vec::new(), None)
2557 };
2558 Ok(ConnectedStore {
2559 event_store,
2560 outbox_store: Some(outbox_store),
2561 namespace_store,
2562 worker_deployment_store,
2563 workloop_store,
2564 visibility_store,
2565 assistant_store,
2566 observability_store: Some(observability_store),
2567 bootstrap_coordinator,
2568 cluster_responder: responder,
2569 cluster_store,
2570 watched_peers,
2571 directory_peers,
2572 self_node_id,
2573 })
2574}
2575
2576fn build_haematite_store(
2591 data_dir: &str,
2592 shard_count: usize,
2593 cluster: Option<crate::config::ClusterConfig>,
2594 node_cache_budget: haematite::NodeCacheBudget,
2595 stage: &crate::control::StageReporter,
2596) -> Result<
2597 (
2598 aion_store_haematite::HaematiteStore,
2599 Option<aion_store_haematite::ClusterResponder>,
2600 ),
2601 ServerError,
2602> {
2603 build_haematite_store_with_hook(
2604 data_dir,
2605 shard_count,
2606 cluster,
2607 node_cache_budget,
2608 stage,
2609 || Ok(()),
2610 )
2611}
2612
2613fn build_haematite_store_with_hook(
2614 data_dir: &str,
2615 shard_count: usize,
2616 cluster: Option<crate::config::ClusterConfig>,
2617 node_cache_budget: haematite::NodeCacheBudget,
2618 stage: &crate::control::StageReporter,
2619 before_backend_touch: impl FnOnce() -> Result<(), std::io::Error>,
2620) -> Result<
2621 (
2622 aion_store_haematite::HaematiteStore,
2623 Option<aion_store_haematite::ClusterResponder>,
2624 ),
2625 ServerError,
2626> {
2627 use aion_store_haematite::{ClusterBootstrap, HaematiteStore};
2628
2629 let private_root = crate::filesystem::ConfinedDir::open_or_create(std::path::Path::new(
2637 data_dir,
2638 ))
2639 .map_err(|error| ServerError::Config {
2640 message: format!("unsafe store.data_dir `{data_dir}`: {error}"),
2641 })?;
2642
2643 stage.report(
2652 crate::control::stage::STAGE_STORE_OPEN,
2653 format!("checking the data root and {shard_count} shard directories under {data_dir}"),
2654 );
2655 for shard in 0..shard_count {
2656 let relative = std::path::PathBuf::from(format!("shard-{shard}"));
2657 private_root
2658 .create_dir_all(&relative)
2659 .map_err(|error| ServerError::Config {
2660 message: format!(
2661 "failed to materialize shard-{shard} under store.data_dir `{data_dir}`: {error}"
2662 ),
2663 })?;
2664 private_root
2665 .ensure_child_dir_private(&relative)
2666 .map_err(|error| private_store_mode_error(data_dir, &error))?;
2667 }
2668
2669 before_backend_touch().map_err(|error| ServerError::Config {
2672 message: format!("store.data_dir pre-open hook failed: {error}"),
2673 })?;
2674
2675 #[cfg(unix)]
2676 let backend_path = private_root
2677 .backend_path()
2678 .map_err(|error| ServerError::Config {
2679 message: format!("failed to resolve held store.data_dir `{data_dir}`: {error}"),
2680 })?;
2681 #[cfg(all(unix, not(any(target_os = "linux", target_os = "android"))))]
2682 crate::filesystem::validate_ambient_backend_ancestors(&backend_path).map_err(|error| {
2683 let (component, reason) = error.into_parts();
2684 ServerError::UnsafeDataRootAncestor {
2685 data_root: backend_path.clone(),
2686 component,
2687 reason,
2688 }
2689 })?;
2690 #[cfg(not(unix))]
2691 let backend_path = std::path::PathBuf::from(data_dir);
2692
2693 let Some(cluster) = cluster else {
2694 let store = if backend_path.join("config.json").exists() {
2695 HaematiteStore::open(
2700 &backend_path,
2701 node_cache_budget,
2702 writer_lock_wait_reporter(stage),
2703 )
2704 .map_err(ServerError::from)?
2705 } else {
2706 HaematiteStore::create_with_shard_count(&backend_path, shard_count, node_cache_budget)
2707 .map_err(ServerError::from)?
2708 };
2709 materialize_narrated(&store, shard_count, stage)?;
2710 let store = store.retain_data_root_capability(private_root);
2711 return Ok((store, None));
2712 };
2713
2714 let boot = ClusterBootstrap {
2715 node_id: cluster.node_id,
2716 bind_address: cluster.bind_address,
2717 members: cluster.members,
2718 peers: cluster
2719 .peers
2720 .into_iter()
2721 .map(|peer| (peer.name, peer.address))
2722 .collect(),
2723 timeout: HAEMATITE_CLUSTER_OP_TIMEOUT,
2724 };
2725 let (store, responder) = HaematiteStore::open_or_create_distributed(
2726 &backend_path,
2727 shard_count,
2728 boot,
2729 node_cache_budget,
2730 )
2731 .map_err(ServerError::from)?;
2732 materialize_narrated(&store, shard_count, stage)?;
2733 let store = store.retain_data_root_capability(private_root);
2734 Ok((store, Some(responder)))
2735}
2736
2737fn writer_lock_wait_reporter(
2745 stage: &crate::control::StageReporter,
2746) -> impl FnMut(&std::path::Path, std::time::Duration) + use<'_> {
2747 move |lock_path, waited| {
2748 stage.report(
2749 crate::control::stage::STAGE_WRITER_LOCK_WAIT,
2750 format!(
2751 "waiting {}s on the store writer lock at {} — another live process \
2752 holds it (a draining predecessor, or another server on this data \
2753 directory); this boot proceeds the moment it is released",
2754 waited.as_secs(),
2755 lock_path.display()
2756 ),
2757 );
2758 }
2759}
2760
2761fn materialize_narrated(
2772 store: &aion_store_haematite::HaematiteStore,
2773 shard_count: usize,
2774 stage: &crate::control::StageReporter,
2775) -> Result<(), ServerError> {
2776 stage.report(
2777 crate::control::stage::STAGE_STORE_OPEN,
2778 format!("haematite store opened; materializing {shard_count} shards"),
2779 );
2780 store
2781 .materialize_all_shards(|shard, materialized, total| {
2782 stage.report(
2783 crate::control::stage::STAGE_WAL_RECOVERY,
2784 format!("materializing shard {shard} ({materialized} of {total})"),
2785 );
2786 })
2787 .map_err(ServerError::from)?;
2788 stage.report(
2789 crate::control::stage::STAGE_WAL_RECOVERY,
2790 format!("all {shard_count} shards materialized"),
2791 );
2792 Ok(())
2793}
2794
2795fn private_store_mode_error(data_dir: &str, error: &std::io::Error) -> ServerError {
2796 ServerError::Config {
2797 message: format!(
2798 "failed to apply private modes under store.data_dir `{data_dir}`: {error}"
2799 ),
2800 }
2801}
2802
2803const HAEMATITE_CLUSTER_OP_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(5);
2805
2806#[cfg(not(feature = "liminal-transport"))]
2815#[derive(Clone, Debug)]
2816struct NullInterventionTransport;
2817
2818#[cfg(not(feature = "liminal-transport"))]
2819#[async_trait::async_trait]
2820impl crate::worker::InterventionTransport for NullInterventionTransport {
2821 async fn push(
2822 &self,
2823 _worker: &crate::worker::WorkerHandle,
2824 _command: aion_core::InterventionCommand,
2825 ) -> Result<aion_core::InterventionOutcome, ServerError> {
2826 Err(ServerError::worker_connection_lost(
2827 "intervention",
2828 "no intervention push transport is compiled in".to_owned(),
2829 ))
2830 }
2831}
2832
2833#[cfg(test)]
2834mod tests {
2835 use std::{net::SocketAddr, time::Duration};
2836
2837 use aion_store::InMemoryStore;
2838
2839 use super::ServerState;
2840 use crate::config::{
2841 AuthConfig, AuthoringConfig, DeployConfig, DevConfig, ListenConfig, MetricsConfig,
2842 NamespaceConfig, NamespaceMode, OpsConsoleAssetSource, OpsConsoleConfig, OutboxConfig,
2843 RuntimeConfig, WebSocketConfig, WorkerConfig,
2844 };
2845
2846 fn runtime_config() -> RuntimeConfig {
2847 RuntimeConfig {
2848 listen: ListenConfig {
2849 grpc: SocketAddr::from(([127, 0, 0, 1], 50051)),
2850 http: SocketAddr::from(([127, 0, 0, 1], 8080)),
2851 },
2852 tls: None,
2853 auth: AuthConfig {
2854 enabled: false,
2855 jwks_url: None,
2856 jwks_refresh_seconds: 300,
2857 },
2858 ops_console: OpsConsoleConfig {
2859 source: OpsConsoleAssetSource::Embedded,
2860 },
2861 namespace: NamespaceConfig {
2862 mode: NamespaceMode::SharedEngine,
2863 },
2864 worker: WorkerConfig {
2865 heartbeat_window: Duration::from_secs(30),
2866 ..WorkerConfig::default()
2867 },
2868 websocket: WebSocketConfig {
2869 outbound_buffer_bound: 32,
2870 event_broadcast_capacity: Some(64),
2871 cluster_broadcast_capacity: Some(64),
2872 },
2873 workflow_packages: Vec::new(),
2874 deploy: DeployConfig::default(),
2875 authoring: AuthoringConfig::default(),
2876 dev: DevConfig::default(),
2877 outbox: OutboxConfig::default(),
2878 observability: crate::config::ObservabilityConfig::with_flush_policy(64, 0),
2879 mcp: crate::config::ResolvedMcpConfig::default(),
2880 assistant: crate::config::ResolvedAssistantConfig::default(),
2881 scheduler_threads: 1,
2882 stop_drain_timeout: Some(std::time::Duration::from_secs(5)),
2883 jit_threshold: None,
2884 query_timeout: Some(Duration::from_secs(10)),
2885 workloop_sweep_interval: Some(Duration::from_millis(50)),
2886 default_namespace: "default".to_owned(),
2887 auto_create: crate::config::AutoCreate::Open,
2888 max_in_flight_activities: crate::config::DEFAULT_MAX_IN_FLIGHT_ACTIVITIES,
2889 drain_timeout: Duration::from_secs(30),
2890 metrics: MetricsConfig { enabled: true },
2891 owned_shards: Vec::new(),
2892 cors_allowed_origins: Vec::new(),
2893 }
2894 }
2895
2896 #[test]
2901 fn an_unruled_flush_policy_refuses_to_build_the_publisher() {
2902 for (mutate, expected_key) in [
2903 (
2904 Box::new(|runtime: &mut RuntimeConfig| {
2905 runtime.observability.max_batch_events = None;
2906 }) as Box<dyn Fn(&mut RuntimeConfig)>,
2907 "observability.max_batch_events",
2908 ),
2909 (
2910 Box::new(|runtime: &mut RuntimeConfig| {
2911 runtime.observability.max_batch_events = Some(0);
2912 }),
2913 "observability.max_batch_events",
2914 ),
2915 (
2916 Box::new(|runtime: &mut RuntimeConfig| {
2917 runtime.observability.max_batch_hold_ms = None;
2918 }),
2919 "observability.max_batch_hold_ms",
2920 ),
2921 ] {
2922 let mut runtime = runtime_config();
2923 mutate(&mut runtime);
2924 let message = super::required_transcript_batch_policy(&runtime)
2925 .err()
2926 .map_or_else(String::new, |error| error.to_string());
2927 assert!(
2928 message.contains(expected_key),
2929 "the refusal must name {expected_key}: {message}"
2930 );
2931 assert!(
2932 message.contains("no default"),
2933 "and must say the key has no default: {message}"
2934 );
2935 }
2936 }
2937
2938 #[test]
2941 fn a_stated_flush_policy_is_carried_through_verbatim() -> Result<(), Box<dyn std::error::Error>>
2942 {
2943 let mut runtime = runtime_config();
2944 runtime.observability = crate::config::ObservabilityConfig::with_flush_policy(32, 0);
2945 let policy = super::required_transcript_batch_policy(&runtime)?;
2946 assert_eq!(policy.max_batch_events.get(), 32);
2947 assert_eq!(policy.max_hold, Duration::ZERO);
2948
2949 runtime.observability = crate::config::ObservabilityConfig::with_flush_policy(8, 250);
2950 let policy = super::required_transcript_batch_policy(&runtime)?;
2951 assert_eq!(policy.max_batch_events.get(), 8);
2952 assert_eq!(policy.max_hold, Duration::from_millis(250));
2953 Ok(())
2954 }
2955
2956 #[test]
2967 fn engine_schema_accepts_every_attribute_the_start_writer_records()
2968 -> Result<(), Box<dyn std::error::Error>> {
2969 let schema = super::server_search_attribute_schema()?;
2970 let recorded = crate::api::handlers::workflows::start_search_attributes(
2971 "tenant-a",
2972 Some("gpu"),
2973 Some("Nightly settlement"),
2974 );
2975
2976 assert!(
2977 recorded.contains_key(crate::namespace::DISPLAY_NAME_ATTRIBUTE),
2978 "the fixture must exercise the display-name attribute, or this test \
2979 cannot see its registration go missing"
2980 );
2981 for (name, value) in &recorded {
2982 schema.validate(name, value).map_err(|error| {
2983 format!(
2984 "the start writer records {name}, but the engine's schema refuses it: {error}"
2985 )
2986 })?;
2987 }
2988 Ok(())
2989 }
2990
2991 #[tokio::test]
2992 async fn builds_state_with_in_memory_store() -> Result<(), Box<dyn std::error::Error>> {
2993 let state =
2994 ServerState::build_with_store(InMemoryStore::default(), runtime_config()).await?;
2995
2996 std::hint::black_box(state.namespace_guard());
2997 std::hint::black_box(state.worker_registry());
2998
2999 Ok(())
3000 }
3001
3002 #[tokio::test]
3010 async fn unserved_queues_surfaces_a_parked_dispatch_and_clears_it()
3011 -> Result<(), Box<dyn std::error::Error>> {
3012 use aion::{ActivityDispatch, ActivityDispatcher as _};
3013 use aion_core::{ActivityId, RunId, WorkflowId};
3014 use std::sync::Arc;
3015
3016 let state =
3017 ServerState::build_with_store(InMemoryStore::default(), runtime_config()).await?;
3018 assert!(
3019 state.unserved_queues()?.is_empty(),
3020 "a calm boot has no unserved queues"
3021 );
3022 assert!(state.queue_declarations().is_installed());
3026 assert_eq!(
3027 state
3028 .queue_declarations()
3029 .declaration_for("nobody-serves-this"),
3030 crate::worker::QueueDeclaration::Unknown
3031 );
3032
3033 let dispatcher = Arc::new(
3034 crate::worker::WorkerActivityDispatcher::new(
3035 state.worker_registry().clone(),
3036 "default",
3037 crate::worker::HeartbeatTracker::new(Duration::from_secs(5)),
3038 )
3039 .with_queue_state(state.queue_service_state().clone())
3040 .with_queue_declarations(state.queue_declarations().clone()),
3041 );
3042 let workflow_id = WorkflowId::new_v4();
3043 let request = ActivityDispatch {
3044 namespace: "default".to_owned(),
3045 task_queue: "nobody-serves-this".to_owned(),
3046 node: None,
3047 workflow_id: workflow_id.clone(),
3048 run_id: RunId::new_v4(),
3049 activity_id: ActivityId::from_sequence_position(0),
3050 name: "greet".to_owned(),
3051 input: "{}".to_owned(),
3052 config: "{}".to_owned(),
3053 attempt: 1,
3054 advisory: false,
3055 labels: std::collections::BTreeMap::new(),
3056 };
3057 let parked = std::thread::spawn(move || dispatcher.dispatch(request));
3058
3059 let mut unserved = Vec::new();
3060 for _ in 0..30 {
3061 unserved = state.unserved_queues()?;
3062 if !unserved.is_empty() {
3063 break;
3064 }
3065 tokio::time::sleep(Duration::from_millis(100)).await;
3066 }
3067 assert_eq!(unserved.len(), 1, "the parked dispatch is not surfaced");
3068 assert_eq!(
3069 unserved[0].reason,
3070 crate::worker::QueueServiceReason::NoLivePollers,
3071 "an empty catalog must not be read as a structural refusal"
3072 );
3073 assert_eq!(unserved[0].key.task_queue, "nobody-serves-this");
3074 assert_eq!(unserved[0].waiting.len(), 1);
3075 assert_eq!(unserved[0].waiting[0].workflow_id, workflow_id);
3076
3077 let (worker_tx, worker_rx) = tokio::sync::mpsc::channel(1);
3079 drop(worker_rx);
3080 let registration = state.worker_registry().register_namespaces(
3081 [String::from("default")],
3082 "nobody-serves-this",
3083 None,
3084 [String::from("greet")].iter(),
3085 worker_tx,
3086 )?;
3087 let outcome = parked.join().map_err(|_| "parked dispatch panicked")?;
3088 assert!(outcome.is_err(), "the released dispatch must resolve");
3089 assert!(
3090 state.unserved_queues()?.is_empty(),
3091 "a resolved dispatch must leave the unserved state"
3092 );
3093 registration.deregister()?;
3094 Ok(())
3095 }
3096
3097 #[tokio::test]
3098 async fn namespace_store_is_reachable_and_functional_after_default_boot()
3099 -> Result<(), Box<dyn std::error::Error>> {
3100 use aion_store::{MintOutcome, NamespaceOrigin};
3101
3102 let state =
3106 ServerState::build_with_store(InMemoryStore::default(), runtime_config()).await?;
3107
3108 let store = state.namespace_store();
3109
3110 let outcome = store
3112 .register_namespace("orders", NamespaceOrigin::WorkerMint)
3113 .await?;
3114 assert_eq!(
3115 outcome,
3116 MintOutcome::Created,
3117 "the first reference to a namespace mints it"
3118 );
3119
3120 let again = store
3122 .register_namespace("orders", NamespaceOrigin::WorkerMint)
3123 .await?;
3124 assert_eq!(
3125 again,
3126 MintOutcome::AlreadyExisted,
3127 "a second reference touches the existing record rather than re-creating it"
3128 );
3129
3130 let fetched = store.get_namespace("orders").await?;
3132 let record = fetched.ok_or("registered namespace must be retrievable via get_namespace")?;
3133 assert_eq!(record.name, "orders");
3134 assert_eq!(record.origin, NamespaceOrigin::WorkerMint);
3135
3136 let listed = store.list_namespaces().await?;
3138 assert!(
3139 listed.iter().any(|record| record.name == "orders"),
3140 "list_namespaces returns the minted namespace"
3141 );
3142
3143 Ok(())
3144 }
3145
3146 #[tokio::test(flavor = "multi_thread")]
3147 async fn connect_store_haematite_round_trips_through_event_store()
3148 -> Result<(), Box<dyn std::error::Error>> {
3149 use aion_core::{ContentType, EventEnvelope, PackageVersion, Payload, RunId, WorkflowId};
3150 use aion_store::WriteToken;
3151 use chrono::Utc;
3152
3153 use crate::config::{StoreBackend, StoreConfig};
3154
3155 let data_dir = crate::test_support::private_tempdir()?;
3156 let connected = super::connect_store(
3160 StoreConfig {
3161 backend: StoreBackend::Haematite,
3162 owned_shards: Vec::new(),
3163 data_dir: Some(data_dir.path().to_string_lossy().into_owned()),
3164 shard_count: 1,
3165 cluster: None,
3166 node_cache_budget: Some(test_node_cache_budget()?),
3167 ..StoreConfig::default()
3168 },
3169 &crate::control::StageReporter::detached(),
3170 )
3171 .await?;
3172 let event_store = connected.event_store;
3173 assert!(
3174 connected.outbox_store.is_some(),
3175 "the haematite backend shares its leaf store as the dispatcher's outbox store"
3176 );
3177 assert!(
3178 connected.bootstrap_coordinator,
3179 "a single-node haematite boot owns all shards and bootstraps the coordinator"
3180 );
3181 assert!(
3182 connected.cluster_responder.is_none(),
3183 "a single-node (no [cluster]) haematite boot has no distributed responder"
3184 );
3185
3186 let workflow_id = WorkflowId::new_v4();
3187 let event = aion_core::Event::WorkflowStarted {
3188 envelope: EventEnvelope {
3189 seq: 1,
3190 recorded_at: Utc::now(),
3191 workflow_id: workflow_id.clone(),
3192 },
3193 workflow_type: String::from("checkout"),
3194 input: Payload::new(ContentType::Json, b"{}".to_vec()),
3195 run_id: RunId::new_v4(),
3196 parent_run_id: None,
3197 parent_workflow_id: None,
3198 package_version: PackageVersion::new("a".repeat(64)),
3199 };
3200 event_store
3201 .append(
3202 WriteToken::recorder(),
3203 &workflow_id,
3204 std::slice::from_ref(&event),
3205 0,
3206 )
3207 .await?;
3208 let history = event_store.read_history(&workflow_id).await?;
3209 assert_eq!(
3210 history.len(),
3211 1,
3212 "an event appended through the server's dyn EventStore reads back"
3213 );
3214 Ok(())
3215 }
3216
3217 fn test_node_cache_budget() -> Result<haematite::NodeCacheBudget, Box<dyn std::error::Error>> {
3223 Ok(haematite::NodeCacheBudget::bytes(1 << 30)?)
3224 }
3225
3226 const BOOT_COST_WORKFLOWS: usize = 800;
3231
3232 const BOOT_OPENS_PER_DIRECTORY: u64 = 8;
3238
3239 #[cfg(unix)]
3256 #[tokio::test]
3257 async fn booting_over_a_store_of_many_object_files_opens_per_shard_not_per_file()
3258 -> Result<(), Box<dyn std::error::Error>> {
3259 let sandbox = crate::test_support::private_tempdir()?;
3260 let data_root = sandbox.path().join("data");
3261 let data_dir = data_root
3262 .to_str()
3263 .ok_or("temporary data path was not UTF-8")?
3264 .to_owned();
3265 let shard_count = 4;
3266
3267 populate_boot_cost_store(&data_root, shard_count).await?;
3268 let object_files = count_regular_files(&data_root)?;
3269 let directories = u64::try_from(shard_count)? + 1;
3270 let budget = BOOT_OPENS_PER_DIRECTORY * directories;
3271 assert!(
3272 object_files > budget * 4,
3273 "the specimen must carry far more object files ({object_files}) than the open \
3274 budget ({budget}), or a walk would fit inside the budget and this pin proves \
3275 nothing"
3276 );
3277
3278 crate::filesystem::reset_capability_opens();
3279 haematite::store::reset_file_opens();
3280 let (store, responder) = tokio::task::spawn_blocking(move || {
3281 super::build_haematite_store_with_hook(
3282 &data_dir,
3283 shard_count,
3284 None,
3285 haematite::NodeCacheBudget::bytes(1 << 30).map_err(|error| error.to_string())?,
3286 &crate::control::StageReporter::detached(),
3287 || Ok(()),
3288 )
3289 .map_err(|error| error.to_string())
3290 })
3291 .await??;
3292 let opens = crate::filesystem::capability_opens();
3293 let haematite_opens = haematite::store::file_opens();
3294 assert!(responder.is_none());
3295 assert_post_boot_node_is_private(&store, &data_root).await?;
3296 drop(store);
3297 println!(
3298 "boot-cost pin: {object_files} object files in {shard_count} shards; {opens} opens \
3299 through the server's capability layer and {haematite_opens} inside haematite during \
3300 the open (budget {budget} each)"
3301 );
3302
3303 assert!(
3304 opens <= budget,
3305 "booting over {object_files} object files in {shard_count} shards issued {opens} \
3306 opens through the server's capability layer; the budget is {budget} \
3307 ({BOOT_OPENS_PER_DIRECTORY} per directory) — the boot is walking the tree again"
3308 );
3309 assert!(
3310 haematite_opens <= budget,
3311 "opening the store over {object_files} object files issued {haematite_opens} opens \
3312 inside haematite; the budget is {budget} — the shard open is doing per-object work"
3313 );
3314 Ok(())
3315 }
3316
3317 #[cfg(unix)]
3321 #[test]
3322 fn a_permissive_shard_directory_we_own_is_tightened_at_boot()
3323 -> Result<(), Box<dyn std::error::Error>> {
3324 use std::os::unix::fs::PermissionsExt as _;
3325
3326 let sandbox = crate::test_support::private_tempdir()?;
3327 let data_root = sandbox.path().join("data");
3328 let loose = data_root.join("shard-2");
3329 std::fs::create_dir_all(&loose)?;
3330 std::fs::set_permissions(&data_root, std::fs::Permissions::from_mode(0o700))?;
3331 std::fs::set_permissions(&loose, std::fs::Permissions::from_mode(0o755))?;
3332 let data_dir = data_root
3333 .to_str()
3334 .ok_or("temporary data path was not UTF-8")?;
3335
3336 let (captured, built) = crate::test_support::CapturedLogs::capture(|| {
3337 super::build_haematite_store_with_hook(
3338 data_dir,
3339 4,
3340 None,
3341 haematite::NodeCacheBudget::bytes(1 << 30).map_err(|error| error.to_string())?,
3342 &crate::control::StageReporter::detached(),
3343 || Ok(()),
3344 )
3345 .map_err(|error| error.to_string())
3346 });
3347 let (store, _) = built?;
3348 drop(store);
3349
3350 assert_eq!(
3351 std::fs::metadata(&loose)?.permissions().mode() & 0o777,
3352 0o700,
3353 "the loose shard directory is tightened by the boot"
3354 );
3355 let text = captured.text()?;
3356 assert!(
3357 text.contains("tightened a sensitive root") && text.contains("shard-2"),
3358 "the repair is logged and names the directory: {text}"
3359 );
3360 Ok(())
3361 }
3362
3363 #[test]
3373 fn the_store_build_narrates_each_shard_by_id_before_opening_it()
3374 -> Result<(), Box<dyn std::error::Error>> {
3375 const SHARDS: usize = 6;
3376 let sandbox = crate::test_support::private_tempdir()?;
3377 let data_root = sandbox.path().join("data");
3378 let data_dir = data_root
3379 .to_str()
3380 .ok_or("temporary data path was not UTF-8")?;
3381
3382 let (captured, built) = crate::test_support::CapturedLogs::capture(|| {
3383 super::build_haematite_store_with_hook(
3384 data_dir,
3385 SHARDS,
3386 None,
3387 haematite::NodeCacheBudget::bytes(1 << 30).map_err(|error| error.to_string())?,
3388 &crate::control::StageReporter::detached(),
3389 || Ok(()),
3390 )
3391 .map_err(|error| error.to_string())
3392 });
3393 let (store, _) = built?;
3394 drop(store);
3395
3396 let text = captured.text()?;
3397 let position = |needle: &str| -> Result<usize, String> {
3398 text.find(needle)
3399 .ok_or_else(|| format!("missing stage line `{needle}` in:\n{text}"))
3400 };
3401 let checking = position(&format!(
3402 "checking the data root and {SHARDS} shard directories"
3403 ))?;
3404 let opened = position(&format!(
3405 "haematite store opened; materializing {SHARDS} shards"
3406 ))?;
3407 let done = position(&format!("all {SHARDS} shards materialized"))?;
3408 assert!(
3409 checking < opened && opened < done,
3410 "stage lines out of order:\n{text}"
3411 );
3412 for shard in 0..SHARDS {
3413 let line = position(&format!("materializing shard {shard} ("))?;
3414 assert!(
3415 opened < line && line < done,
3416 "shard {shard}'s start line must sit between open and done:\n{text}"
3417 );
3418 }
3419 assert!(
3420 !text.contains(&format!("materializing shard {SHARDS} (")),
3421 "a line names a shard id; ids run 0..{SHARDS}, never the count"
3422 );
3423 Ok(())
3424 }
3425
3426 async fn populate_boot_cost_store(
3429 data_root: &std::path::Path,
3430 shard_count: usize,
3431 ) -> Result<(), Box<dyn std::error::Error>> {
3432 use aion_core::{ContentType, EventEnvelope, PackageVersion, Payload, RunId, WorkflowId};
3433 use aion_store::{WritableEventStore as _, WriteToken};
3434 use chrono::Utc;
3435
3436 let store = aion_store_haematite::HaematiteStore::create_with_shard_count(
3437 data_root,
3438 shard_count,
3439 test_node_cache_budget()?,
3440 )?;
3441 for index in 0..BOOT_COST_WORKFLOWS {
3442 let workflow_id = WorkflowId::new_v4();
3443 let event = aion_core::Event::WorkflowStarted {
3444 envelope: EventEnvelope {
3445 seq: 1,
3446 recorded_at: Utc::now(),
3447 workflow_id: workflow_id.clone(),
3448 },
3449 workflow_type: format!("boot-cost-{index}"),
3450 input: Payload::new(ContentType::Json, b"{}".to_vec()),
3451 run_id: RunId::new_v4(),
3452 parent_run_id: None,
3453 parent_workflow_id: None,
3454 package_version: PackageVersion::new("a".repeat(64)),
3455 };
3456 store
3457 .append(
3458 WriteToken::recorder(),
3459 &workflow_id,
3460 std::slice::from_ref(&event),
3461 0,
3462 )
3463 .await?;
3464 }
3465 Ok(())
3466 }
3467
3468 async fn assert_post_boot_node_is_private(
3476 store: &aion_store_haematite::HaematiteStore,
3477 data_root: &std::path::Path,
3478 ) -> Result<(), Box<dyn std::error::Error>> {
3479 use std::os::unix::fs::PermissionsExt as _;
3480
3481 use aion_core::{ContentType, EventEnvelope, PackageVersion, Payload, RunId, WorkflowId};
3482 use aion_store::{WritableEventStore as _, WriteToken};
3483 use chrono::Utc;
3484
3485 let workflow_id = WorkflowId::new_v4();
3486 let event = aion_core::Event::WorkflowStarted {
3487 envelope: EventEnvelope {
3488 seq: 1,
3489 recorded_at: Utc::now(),
3490 workflow_id: workflow_id.clone(),
3491 },
3492 workflow_type: String::from("boot-cost-after-boot"),
3493 input: Payload::new(ContentType::Json, b"{}".to_vec()),
3494 run_id: RunId::new_v4(),
3495 parent_run_id: None,
3496 parent_workflow_id: None,
3497 package_version: PackageVersion::new("a".repeat(64)),
3498 };
3499 store
3500 .append(
3501 WriteToken::recorder(),
3502 &workflow_id,
3503 std::slice::from_ref(&event),
3504 0,
3505 )
3506 .await?;
3507 let (newest, parent) = newest_regular_file(data_root)?
3508 .ok_or("the post-boot append must have written a node file")?;
3509 assert_eq!(
3510 std::fs::metadata(&newest)?.permissions().mode() & 0o777,
3511 0o600,
3512 "a node written after boot is private by creation: {}",
3513 newest.display()
3514 );
3515 assert_eq!(
3516 std::fs::metadata(&parent)?.permissions().mode() & 0o777,
3517 0o700,
3518 "its prefix directory is private by creation: {}",
3519 parent.display()
3520 );
3521 Ok(())
3522 }
3523
3524 fn newest_regular_file(
3526 root: &std::path::Path,
3527 ) -> std::io::Result<Option<(std::path::PathBuf, std::path::PathBuf)>> {
3528 let mut newest: Option<(std::time::SystemTime, std::path::PathBuf)> = None;
3529 let mut pending = vec![root.to_path_buf()];
3530 while let Some(dir) = pending.pop() {
3531 for entry in std::fs::read_dir(&dir)? {
3532 let entry = entry?;
3533 let file_type = entry.file_type()?;
3534 if file_type.is_dir() {
3535 pending.push(entry.path());
3536 } else if file_type.is_file() {
3537 let modified = entry.metadata()?.modified()?;
3538 if newest.as_ref().is_none_or(|(when, _)| modified > *when) {
3539 newest = Some((modified, entry.path()));
3540 }
3541 }
3542 }
3543 }
3544 Ok(newest.and_then(|(_, path)| {
3545 let parent = path.parent()?.to_path_buf();
3546 Some((path, parent))
3547 }))
3548 }
3549
3550 fn count_regular_files(root: &std::path::Path) -> std::io::Result<u64> {
3551 let mut total = 0;
3552 for entry in std::fs::read_dir(root)? {
3553 let entry = entry?;
3554 let file_type = entry.file_type()?;
3555 if file_type.is_dir() {
3556 total += count_regular_files(&entry.path())?;
3557 } else if file_type.is_file() {
3558 total += 1;
3559 }
3560 }
3561 Ok(total)
3562 }
3563
3564 #[cfg(unix)]
3565 #[test]
3566 fn haematite_root_swap_before_first_backend_touch_cannot_redirect_writes()
3567 -> Result<(), Box<dyn std::error::Error>> {
3568 use std::os::unix::fs::symlink;
3569
3570 let sandbox = crate::test_support::private_tempdir()?;
3571 let configured_root = sandbox.path().join("data");
3572 let held_root = sandbox.path().join("held-data");
3573 let outside = sandbox.path().join("outside");
3574 std::fs::create_dir(&outside)?;
3575 let configured = configured_root
3576 .to_str()
3577 .ok_or("temporary data path was not UTF-8")?;
3578
3579 let (store, responder) = super::build_haematite_store_with_hook(
3580 configured,
3581 4,
3582 None,
3583 test_node_cache_budget()?,
3584 &crate::control::StageReporter::detached(),
3585 || {
3586 std::fs::rename(&configured_root, &held_root)?;
3591 symlink(&outside, &configured_root)?;
3592 Ok(())
3593 },
3594 )?;
3595 assert!(responder.is_none());
3596
3597 let outside_entries = std::fs::read_dir(&outside)?.collect::<Result<Vec<_>, _>>()?;
3598 assert!(
3599 outside_entries.is_empty(),
3600 "Haematite followed the replaced ambient root and wrote outside"
3601 );
3602 assert!(held_root.join("config.json").is_file());
3603 for shard in 0..4 {
3604 let shard_path = held_root.join(format!("shard-{shard}"));
3605 assert!(shard_path.is_dir(), "shard {shard} was not materialized");
3606 assert!(
3607 std::fs::read_dir(&shard_path)?
3608 .next()
3609 .transpose()?
3610 .is_some(),
3611 "shard {shard} did not run Haematite's materialization path"
3612 );
3613 }
3614
3615 drop(store);
3616 Ok(())
3617 }
3618
3619 #[cfg(any(target_os = "linux", target_os = "android"))]
3620 #[tokio::test]
3621 async fn proc_fd_backend_path_survives_a_post_startup_root_swap()
3622 -> Result<(), Box<dyn std::error::Error>> {
3623 use std::os::unix::fs::symlink;
3624
3625 use aion_core::{ContentType, EventEnvelope, PackageVersion, Payload, RunId, WorkflowId};
3626 use aion_store::{WritableEventStore as _, WriteToken};
3627 use chrono::Utc;
3628
3629 let sandbox = crate::test_support::private_tempdir()?;
3630 let configured_root = sandbox.path().join("data");
3631 let held_root = sandbox.path().join("held-data");
3632 let capture = sandbox.path().join("capture");
3633 std::fs::create_dir(&capture)?;
3634 let configured = configured_root
3635 .to_str()
3636 .ok_or("temporary data path was not UTF-8")?;
3637
3638 let (store, responder) = super::build_haematite_store(
3639 configured,
3640 4,
3641 None,
3642 test_node_cache_budget()?,
3643 &crate::control::StageReporter::detached(),
3644 )?;
3645 assert!(responder.is_none());
3646 std::fs::rename(&configured_root, &held_root)?;
3647 symlink(&capture, &configured_root)?;
3648
3649 let workflow_id = WorkflowId::new_v4();
3650 let event = aion_core::Event::WorkflowStarted {
3651 envelope: EventEnvelope {
3652 seq: 1,
3653 recorded_at: Utc::now(),
3654 workflow_id: workflow_id.clone(),
3655 },
3656 workflow_type: String::from("post-startup-root-swap"),
3657 input: Payload::new(ContentType::Json, b"{}".to_vec()),
3658 run_id: RunId::new_v4(),
3659 parent_run_id: None,
3660 parent_workflow_id: None,
3661 package_version: PackageVersion::new("a".repeat(64)),
3662 };
3663 store
3664 .append(
3665 WriteToken::recorder(),
3666 &workflow_id,
3667 std::slice::from_ref(&event),
3668 0,
3669 )
3670 .await?;
3671
3672 let captured = std::fs::read_dir(&capture)?.collect::<Result<Vec<_>, _>>()?;
3673 assert!(
3674 captured.is_empty(),
3675 "post-startup append followed the replacement symlink into capture"
3676 );
3677 assert!(held_root.join("config.json").is_file());
3678 drop(store);
3679 Ok(())
3680 }
3681
3682 #[cfg(all(unix, not(any(target_os = "linux", target_os = "android"))))]
3683 #[test]
3684 fn path_ambient_haematite_refuses_group_or_world_writable_ancestors()
3685 -> Result<(), Box<dyn std::error::Error>> {
3686 use std::os::unix::fs::PermissionsExt as _;
3687
3688 let sandbox = crate::test_support::private_tempdir()?;
3689 std::fs::set_permissions(sandbox.path(), std::fs::Permissions::from_mode(0o700))?;
3690
3691 for mode in [0o770, 0o1777] {
3692 let shared = sandbox.path().join(format!("shared-{mode:o}"));
3693 let data_root = shared.join("data");
3694 std::fs::create_dir(&shared)?;
3695 std::fs::set_permissions(&shared, std::fs::Permissions::from_mode(mode))?;
3696 std::fs::create_dir(&data_root)?;
3697 std::fs::set_permissions(&data_root, std::fs::Permissions::from_mode(0o700))?;
3698 let configured = data_root
3699 .to_str()
3700 .ok_or("temporary data path was not UTF-8")?;
3701
3702 let Err(error) = super::build_haematite_store(
3703 configured,
3704 4,
3705 None,
3706 test_node_cache_budget()?,
3707 &crate::control::StageReporter::detached(),
3708 ) else {
3709 return Err(format!("mode {mode:04o} ancestor was accepted").into());
3710 };
3711 let message = error.to_string();
3712 let crate::ServerError::UnsafeDataRootAncestor {
3713 data_root: resolved_root,
3714 component,
3715 reason,
3716 } = error
3717 else {
3718 return Err(format!("expected typed unsafe-ancestor error, got {message}").into());
3719 };
3720 assert_eq!(resolved_root, std::fs::canonicalize(&data_root)?);
3721 assert_eq!(component, std::fs::canonicalize(&shared)?);
3722 assert!(
3723 reason.contains(&format!("mode {mode:04o}")),
3724 "unexpected reason: {reason}"
3725 );
3726 if mode & 0o1000 != 0 {
3727 assert!(reason.contains("sticky bit is not accepted"));
3728 }
3729 assert!(message.contains("private Aion home"));
3730 assert!(
3731 !data_root.join("config.json").exists(),
3732 "Haematite touched its ambient path before the refusal"
3733 );
3734 }
3735 Ok(())
3736 }
3737
3738 #[cfg(target_os = "macos")]
3739 #[test]
3740 fn path_ambient_haematite_refuses_mutating_allow_acl_ancestor()
3741 -> Result<(), Box<dyn std::error::Error>> {
3742 use std::os::unix::fs::PermissionsExt as _;
3743
3744 let sandbox = crate::test_support::private_tempdir()?;
3745 std::fs::set_permissions(sandbox.path(), std::fs::Permissions::from_mode(0o700))?;
3746 let shared = sandbox.path().join("acl-shared");
3747 let data_root = shared.join("data");
3748 std::fs::create_dir(&shared)?;
3749 std::fs::set_permissions(&shared, std::fs::Permissions::from_mode(0o700))?;
3750 let acl = "everyone allow list,search,add_file,add_subdirectory,delete_child";
3751 let status = std::process::Command::new("chmod")
3752 .arg("+a")
3753 .arg(acl)
3754 .arg(&shared)
3755 .status()?;
3756 assert!(status.success(), "failed to install Darwin regression ACL");
3757 let configured = data_root
3758 .to_str()
3759 .ok_or("temporary data path was not UTF-8")?;
3760
3761 let result = super::build_haematite_store(
3762 configured,
3763 4,
3764 None,
3765 test_node_cache_budget()?,
3766 &crate::control::StageReporter::detached(),
3767 );
3768 let cleanup = std::process::Command::new("chmod")
3769 .arg("-RN")
3770 .arg(&shared)
3771 .status()?;
3772 assert!(cleanup.success(), "failed to clean Darwin regression ACL");
3773
3774 let Err(error) = result else {
3775 return Err("mutating non-euid allow ACL ancestor was accepted".into());
3776 };
3777 let message = error.to_string();
3778 let crate::ServerError::UnsafeDataRootAncestor {
3779 component, reason, ..
3780 } = error
3781 else {
3782 return Err(format!("expected typed unsafe-ancestor error, got {message}").into());
3783 };
3784 assert_eq!(component, std::fs::canonicalize(&shared)?);
3785 assert!(
3786 reason.contains("allow"),
3787 "reason did not name the ACE: {reason}"
3788 );
3789 assert!(
3790 reason.contains("everyone"),
3791 "reason did not name the ACE principal: {reason}"
3792 );
3793 assert!(
3794 !data_root.join("config.json").exists(),
3795 "Haematite touched its ambient path before the ACL refusal"
3796 );
3797 Ok(())
3798 }
3799
3800 #[cfg(target_os = "macos")]
3801 #[test]
3802 fn path_ambient_haematite_accepts_the_euid_uuid_allow_ace()
3803 -> Result<(), Box<dyn std::error::Error>> {
3804 use std::os::unix::fs::PermissionsExt as _;
3805
3806 use exacl::{AclEntry, AclOption, Perm};
3807
3808 let sandbox = crate::test_support::private_tempdir()?;
3809 std::fs::set_permissions(sandbox.path(), std::fs::Permissions::from_mode(0o700))?;
3810 let private_parent = sandbox.path().join("euid-uuid-allow");
3811 let data_root = private_parent.join("data");
3812 std::fs::create_dir(&private_parent)?;
3813 std::fs::set_permissions(&private_parent, std::fs::Permissions::from_mode(0o700))?;
3814
3815 let server_uid = rustix::process::geteuid().as_raw();
3816 let ace_qualifier = crate::filesystem::darwin_user_uuid_for_test(server_uid)?;
3817 let entry = AclEntry::allow_user(
3818 &ace_qualifier.to_string(),
3819 Perm::EXECUTE | Perm::WRITE | Perm::APPEND | Perm::DELETE_CHILD,
3820 None,
3821 );
3822 exacl::setfacl(
3823 &[private_parent.as_path()],
3824 &[entry],
3825 AclOption::SYMLINK_ACL,
3826 )?;
3827 let configured = data_root
3828 .to_str()
3829 .ok_or("temporary data path was not UTF-8")?;
3830
3831 let result = super::build_haematite_store(
3832 configured,
3833 4,
3834 None,
3835 test_node_cache_budget()?,
3836 &crate::control::StageReporter::detached(),
3837 );
3838 let cleanup = std::process::Command::new("chmod")
3839 .arg("-RN")
3840 .arg(&private_parent)
3841 .status()?;
3842 assert!(cleanup.success(), "failed to clean euid UUID allow ACL");
3843
3844 let (store, responder) = result?;
3845 assert!(responder.is_none());
3846 assert!(data_root.join("config.json").is_file());
3847 drop(store);
3848 Ok(())
3849 }
3850
3851 #[cfg(target_os = "macos")]
3852 #[test]
3853 fn path_ambient_haematite_refuses_a_non_euid_user_uuid_allow_ace()
3854 -> Result<(), Box<dyn std::error::Error>> {
3855 use std::os::unix::fs::PermissionsExt as _;
3856
3857 use exacl::{AclEntry, AclOption, Perm};
3858
3859 let sandbox = crate::test_support::private_tempdir()?;
3860 std::fs::set_permissions(sandbox.path(), std::fs::Permissions::from_mode(0o700))?;
3861 let shared = sandbox.path().join("non-euid-uuid-allow");
3862 let data_root = shared.join("data");
3863 std::fs::create_dir(&shared)?;
3864 std::fs::set_permissions(&shared, std::fs::Permissions::from_mode(0o700))?;
3865
3866 let server_uid = rustix::process::geteuid().as_raw();
3867 let foreign_uid = u32::from(server_uid == 0);
3868 let foreign_qualifier = crate::filesystem::darwin_user_uuid_for_test(foreign_uid)?;
3869 let entry = AclEntry::allow_user(
3870 &foreign_qualifier.to_string(),
3871 Perm::EXECUTE | Perm::WRITE | Perm::APPEND | Perm::DELETE_CHILD,
3872 None,
3873 );
3874 exacl::setfacl(&[shared.as_path()], &[entry], AclOption::SYMLINK_ACL)?;
3875 let configured = data_root
3876 .to_str()
3877 .ok_or("temporary data path was not UTF-8")?;
3878
3879 let result = super::build_haematite_store(
3880 configured,
3881 4,
3882 None,
3883 test_node_cache_budget()?,
3884 &crate::control::StageReporter::detached(),
3885 );
3886 let cleanup = std::process::Command::new("chmod")
3887 .arg("-RN")
3888 .arg(&shared)
3889 .status()?;
3890 assert!(cleanup.success(), "failed to clean non-euid UUID allow ACL");
3891
3892 let Err(error) = result else {
3893 return Err("mutating non-euid user UUID allow ACE was accepted".into());
3894 };
3895 let message = error.to_string();
3896 let crate::ServerError::UnsafeDataRootAncestor {
3897 component, reason, ..
3898 } = error
3899 else {
3900 return Err(format!("expected typed unsafe-ancestor error, got {message}").into());
3901 };
3902 assert_eq!(component, std::fs::canonicalize(&shared)?);
3903 assert!(
3904 reason.contains("allow") && reason.contains(&format!("server euid {server_uid}")),
3905 "reason did not name the rejected ACE: {reason}"
3906 );
3907 assert!(
3908 !data_root.join("config.json").exists(),
3909 "Haematite touched its ambient path before the UUID ACL refusal"
3910 );
3911 Ok(())
3912 }
3913
3914 #[cfg(target_os = "macos")]
3915 #[test]
3916 fn path_ambient_haematite_accepts_a_deny_only_acl_ancestor()
3917 -> Result<(), Box<dyn std::error::Error>> {
3918 use std::os::unix::fs::PermissionsExt as _;
3919
3920 let sandbox = crate::test_support::private_tempdir()?;
3921 std::fs::set_permissions(sandbox.path(), std::fs::Permissions::from_mode(0o700))?;
3922 let private_parent = sandbox.path().join("deny-only");
3923 let data_root = private_parent.join("data");
3924 std::fs::create_dir(&private_parent)?;
3925 std::fs::set_permissions(&private_parent, std::fs::Permissions::from_mode(0o700))?;
3926 let status = std::process::Command::new("chmod")
3927 .arg("+a")
3928 .arg("everyone deny delete")
3929 .arg(&private_parent)
3930 .status()?;
3931 assert!(status.success(), "failed to install Darwin deny-only ACL");
3932 let configured = data_root
3933 .to_str()
3934 .ok_or("temporary data path was not UTF-8")?;
3935
3936 let result = super::build_haematite_store(
3937 configured,
3938 4,
3939 None,
3940 test_node_cache_budget()?,
3941 &crate::control::StageReporter::detached(),
3942 );
3943 let cleanup = std::process::Command::new("chmod")
3944 .arg("-RN")
3945 .arg(&private_parent)
3946 .status()?;
3947 assert!(cleanup.success(), "failed to clean Darwin deny-only ACL");
3948
3949 let (store, responder) = result?;
3950 assert!(responder.is_none());
3951 assert!(data_root.join("config.json").is_file());
3952 drop(store);
3953 Ok(())
3954 }
3955
3956 #[cfg(target_os = "macos")]
3957 #[test]
3958 fn path_ambient_haematite_accepts_the_stock_home_acl_chain()
3959 -> Result<(), Box<dyn std::error::Error>> {
3960 use std::os::unix::fs::PermissionsExt as _;
3961 use users::os::unix::UserExt as _;
3962
3963 let effective_uid = rustix::process::geteuid().as_raw();
3964 let effective_user = users::get_user_by_uid(effective_uid)
3965 .ok_or_else(|| format!("server euid {effective_uid} has no account record"))?;
3966 let sandbox = tempfile::Builder::new()
3967 .prefix(".aion-acl-home-proof-")
3968 .tempdir_in(effective_user.home_dir())?;
3969 std::fs::set_permissions(sandbox.path(), std::fs::Permissions::from_mode(0o700))?;
3970 let data_root = sandbox.path().join("data");
3971 let configured = data_root
3972 .to_str()
3973 .ok_or("temporary data path was not UTF-8")?;
3974
3975 let (store, responder) = super::build_haematite_store(
3976 configured,
3977 4,
3978 None,
3979 test_node_cache_budget()?,
3980 &crate::control::StageReporter::detached(),
3981 )?;
3982 assert!(responder.is_none());
3983 assert!(data_root.join("config.json").is_file());
3984 drop(store);
3985 Ok(())
3986 }
3987
3988 #[cfg(all(unix, not(any(target_os = "linux", target_os = "android"))))]
3989 #[test]
3990 fn path_ambient_haematite_accepts_an_owner_controlled_chain()
3991 -> Result<(), Box<dyn std::error::Error>> {
3992 use std::os::unix::fs::PermissionsExt as _;
3993
3994 let sandbox = crate::test_support::private_tempdir()?;
3995 std::fs::set_permissions(sandbox.path(), std::fs::Permissions::from_mode(0o700))?;
3996 let private_parent = sandbox.path().join("private");
3997 let data_root = private_parent.join("data");
3998 std::fs::create_dir(&private_parent)?;
3999 std::fs::set_permissions(&private_parent, std::fs::Permissions::from_mode(0o700))?;
4000 let configured = data_root
4001 .to_str()
4002 .ok_or("temporary data path was not UTF-8")?;
4003
4004 let (store, responder) = super::build_haematite_store(
4005 configured,
4006 4,
4007 None,
4008 test_node_cache_budget()?,
4009 &crate::control::StageReporter::detached(),
4010 )?;
4011 assert!(responder.is_none());
4012 assert!(data_root.join("config.json").is_file());
4013 for shard in 0..4 {
4014 assert!(data_root.join(format!("shard-{shard}")).is_dir());
4015 }
4016 drop(store);
4017 Ok(())
4018 }
4019
4020 #[tokio::test]
4021 async fn connect_store_memory_backend_exposes_no_outbox_store()
4022 -> Result<(), Box<dyn std::error::Error>> {
4023 use crate::config::{StoreBackend, StoreConfig};
4024
4025 let connected = super::connect_store(
4028 StoreConfig {
4029 backend: StoreBackend::Memory,
4030 owned_shards: Vec::new(),
4031 data_dir: None,
4032 shard_count: 1,
4033 cluster: None,
4034 node_cache_budget: None,
4035 ..StoreConfig::default()
4036 },
4037 &crate::control::StageReporter::detached(),
4038 )
4039 .await?;
4040 assert!(
4041 connected.outbox_store.is_none(),
4042 "the in-memory backend exposes no outbox store"
4043 );
4044 Ok(())
4045 }
4046
4047 #[tokio::test]
4053 async fn haematite_boot_refuses_without_a_node_cache_budget()
4054 -> Result<(), Box<dyn std::error::Error>> {
4055 use crate::ServerError;
4056 use crate::config::{StoreBackend, StoreConfig};
4057
4058 let sandbox = crate::test_support::private_tempdir()?;
4059 let data_dir = sandbox.path().join("data");
4060 let error = super::connect_haematite_store(
4061 StoreConfig {
4062 backend: StoreBackend::Haematite,
4063 data_dir: Some(
4064 data_dir
4065 .to_str()
4066 .ok_or("temporary data path was not UTF-8")?
4067 .to_owned(),
4068 ),
4069 shard_count: 4,
4070 ..StoreConfig::default()
4071 },
4072 &crate::control::StageReporter::detached(),
4073 )
4074 .await
4075 .err()
4076 .ok_or("the haematite boot path must refuse a store config with no node_cache_budget")?;
4077 let ServerError::Config { message } = error else {
4078 return Err(format!("expected a config refusal, got {error:?}").into());
4079 };
4080 assert!(
4081 message.contains("store.node_cache_budget"),
4082 "the refusal must name the missing key, got: {message}"
4083 );
4084 assert!(
4085 message.contains("AION_STORE_NODE_CACHE_BUDGET"),
4086 "the refusal must name the environment override, got: {message}"
4087 );
4088 Ok(())
4089 }
4090
4091 #[tokio::test]
4096 async fn configured_node_cache_budget_reaches_the_created_database()
4097 -> Result<(), Box<dyn std::error::Error>> {
4098 use crate::config::{StoreBackend, StoreConfig};
4099
4100 const ONE_GIB: usize = 1 << 30;
4101
4102 let sandbox = crate::test_support::private_tempdir()?;
4103 let data_dir = sandbox.path().join("data");
4104 let connected = super::connect_haematite_store(
4105 StoreConfig {
4106 backend: StoreBackend::Haematite,
4107 data_dir: Some(
4108 data_dir
4109 .to_str()
4110 .ok_or("temporary data path was not UTF-8")?
4111 .to_owned(),
4112 ),
4113 shard_count: 4,
4114 node_cache_budget: Some(haematite::NodeCacheBudget::bytes(ONE_GIB)?),
4115 ..StoreConfig::default()
4116 },
4117 &crate::control::StageReporter::detached(),
4118 )
4119 .await?;
4120 drop(connected);
4121
4122 let recorded: serde_json::Value =
4123 serde_json::from_slice(&std::fs::read(data_dir.join("config.json"))?)?;
4124 assert_eq!(
4125 recorded.get("node_cache_budget"),
4126 Some(&serde_json::json!({ "bytes": ONE_GIB })),
4127 "the operator's budget must be the one haematite created the database with"
4128 );
4129 Ok(())
4130 }
4131
4132 #[tokio::test]
4133 async fn state_build_fails_without_event_broadcast_capacity()
4134 -> Result<(), Box<dyn std::error::Error>> {
4135 let mut runtime = runtime_config();
4136 runtime.websocket.event_broadcast_capacity = None;
4137
4138 let error = ServerState::build_with_store(InMemoryStore::default(), runtime)
4139 .await
4140 .err()
4141 .ok_or("state build must fail when event streaming is unsized")?;
4142
4143 assert!(error.is_config(), "expected a config error, got {error}");
4144 assert!(
4145 error
4146 .to_string()
4147 .contains("websocket.event_broadcast_capacity"),
4148 "error must name the missing key: {error}"
4149 );
4150 Ok(())
4151 }
4152
4153 #[tokio::test]
4154 async fn state_build_fails_without_query_timeout() -> Result<(), Box<dyn std::error::Error>> {
4155 let mut runtime = runtime_config();
4156 runtime.query_timeout = None;
4157
4158 let error = ServerState::build_with_store(InMemoryStore::default(), runtime)
4159 .await
4160 .err()
4161 .ok_or("state build must fail when the query reply deadline is unset")?;
4162
4163 assert!(error.is_config(), "expected a config error, got {error}");
4164 assert!(
4165 error.to_string().contains("runtime.query_timeout_ms"),
4166 "error must name the missing key: {error}"
4167 );
4168 assert!(
4169 error.to_string().contains("AION_RUNTIME_QUERY_TIMEOUT_MS"),
4170 "error must name the environment override: {error}"
4171 );
4172 Ok(())
4173 }
4174
4175 #[tokio::test]
4176 async fn state_build_fails_with_zero_query_timeout() -> Result<(), Box<dyn std::error::Error>> {
4177 let mut runtime = runtime_config();
4178 runtime.query_timeout = Some(Duration::ZERO);
4179
4180 let error = ServerState::build_with_store(InMemoryStore::default(), runtime)
4181 .await
4182 .err()
4183 .ok_or("state build must fail when the query reply deadline is zero")?;
4184
4185 assert!(error.is_config(), "expected a config error, got {error}");
4186 assert!(
4187 error.to_string().contains("runtime.query_timeout_ms"),
4188 "error must name the zero-valued key: {error}"
4189 );
4190 Ok(())
4191 }
4192
4193 #[tokio::test(flavor = "multi_thread")]
4216 async fn a_completed_check_through_the_built_dispatcher_lands_in_the_returned_slot()
4217 -> Result<(), Box<dyn std::error::Error>> {
4218 use std::collections::{BTreeMap, VecDeque};
4219 use std::sync::{Arc, Mutex};
4220
4221 use aion::ActivityDispatch;
4222 use aion_core::{ActivityId, RunId, WorkflowId};
4223 use aion_package::ActionBodyContract;
4224
4225 use crate::update_check::document::{FETCH_ACTION, FETCH_COMMAND, UPDATE_CHECK_QUEUE};
4226 use crate::worker::{DeclaredBodies, DeclaredBodyLookup, DispatchingRun};
4227
4228 struct SequencedBodies {
4231 replies: Mutex<VecDeque<DeclaredBodyLookup>>,
4232 }
4233
4234 impl DeclaredBodies for SequencedBodies {
4235 fn body_for(
4236 &self,
4237 _task_queue: &str,
4238 _action: &str,
4239 _run: DispatchingRun<'_>,
4240 ) -> DeclaredBodyLookup {
4241 let mut replies = match self.replies.lock() {
4242 Ok(replies) => replies,
4243 Err(poisoned) => poisoned.into_inner(),
4244 };
4245 replies.pop_front().unwrap_or(DeclaredBodyLookup::None)
4246 }
4247 }
4248
4249 let runtime = runtime_config();
4250 let cluster_publisher = crate::cluster_publisher::ClusterEventPublisher::new(
4251 ServerState::FALLBACK_CLUSTER_BROADCAST_CAPACITY,
4252 );
4253 let namespace_store: Arc<dyn aion_store::NamespaceStore> =
4254 Arc::new(InMemoryStore::default());
4255 let worker_deployment_store: Arc<dyn aion_store::WorkerDeploymentStore> =
4256 Arc::new(InMemoryStore::default());
4257 let metrics = crate::observability::Metrics::new()?;
4258 let seams = super::build_worker_seams(
4259 &runtime,
4260 &cluster_publisher,
4261 &metrics,
4262 &namespace_store,
4263 &worker_deployment_store,
4264 None,
4265 );
4266
4267 let fixture = concat!(
4270 env!("CARGO_MANIFEST_DIR"),
4271 "/src/update_check/fixtures/aion-cli-index.jsonl"
4272 );
4273 seams.declared_bodies.install(Arc::new(SequencedBodies {
4274 replies: Mutex::new(VecDeque::from([
4275 DeclaredBodyLookup::Declared(ActionBodyContract::Run {
4276 command: FETCH_COMMAND.to_owned(),
4277 }),
4278 DeclaredBodyLookup::Declared(ActionBodyContract::Run {
4279 command: format!("cat {fixture}"),
4280 }),
4281 ])),
4282 }));
4283
4284 let transcript = crate::activity_publisher::ActivityEventPublisher::new(
4285 Arc::new(aion_store::InMemoryObservabilityStore::default()),
4286 ServerState::FALLBACK_CLUSTER_BROADCAST_CAPACITY,
4287 crate::activity_publisher::TranscriptBatchPolicy::UNBATCHED,
4288 );
4289 let (dispatcher, _mock_registry, _attempt_owners, _workspace_root, update_status) =
4290 super::build_decorated_dispatcher(&runtime, &seams, transcript);
4291
4292 assert_eq!(
4293 update_status.last(),
4294 None,
4295 "the returned slot must start honestly empty"
4296 );
4297
4298 let dispatch = ActivityDispatch {
4299 namespace: "default".to_owned(),
4300 task_queue: UPDATE_CHECK_QUEUE.to_owned(),
4301 node: None,
4302 workflow_id: WorkflowId::new_v4(),
4303 run_id: RunId::new_v4(),
4304 activity_id: ActivityId::from_sequence_position(1),
4305 name: FETCH_ACTION.to_owned(),
4306 input: "{}".to_owned(),
4307 config: "{}".to_owned(),
4308 attempt: 1,
4309 labels: BTreeMap::new(),
4310 advisory: false,
4311 };
4312 let handle = tokio::task::spawn_blocking(move || dispatcher.dispatch(dispatch));
4313 let encoded = handle
4314 .await?
4315 .map_err(|error| format!("the check dispatch failed: {error}"))?;
4316 let outcome: serde_json::Value = serde_json::from_str(&encoded)?;
4317 assert_eq!(outcome["exit_code"], 0, "the local stand-in command ran");
4318
4319 let recorded = update_status.last().ok_or(
4320 "the completed check must land in the RETURNED slot — the one the boot path \
4321 stores and /update-status serves; an empty slot here is the disconnected-\
4322 producer mis-wire the r1 review proved unmeasured",
4323 )?;
4324 assert_eq!(recorded.latest_known, "0.13.7");
4325 Ok(())
4326 }
4327}