use std::{path::PathBuf, sync::Arc};
use aion::{
ActivityDispatcher, EngineBuilder, RuntimeHandle, SignalRouter, signal::ConcreteSignalRouter,
};
use aion_store::{EventStore, NamespaceStore, OutboxStore, WorkerDeploymentStore};
use crate::dev_ui::{ActivityMockRegistry, DevMockingDispatcher};
#[cfg(feature = "auth")]
use crate::auth::JwksCache;
use crate::{
config::{RuntimeConfig, ServerConfig, StoreBackend, StoreConfig},
error::ServerError,
namespace::{NamespaceGuard, NamespaceMinter, resolver::NamespaceResolver},
observability::{
Metrics, health::HealthState, instrumented_store::InstrumentedEventStore,
metrics::MetricsError,
},
shutdown::DrainState,
worker::{
ConnectedWorkerRegistry, HeartbeatTracker, PendingActivities, WorkerActivityDispatcher,
supervisor::WorkerSupervisor,
},
};
fn new_supervisor(
store: &Arc<dyn WorkerDeploymentStore>,
publisher: &crate::cluster_publisher::ClusterEventPublisher,
) -> Arc<WorkerSupervisor> {
Arc::new(WorkerSupervisor::new(Arc::clone(store), publisher.clone()))
}
#[derive(Clone)]
pub struct ServerState {
inner: Arc<ServerStateInner>,
}
struct ServerStateInner {
namespace_guard: NamespaceGuard,
runtime: RuntimeConfig,
worker_registry: ConnectedWorkerRegistry,
pending_activities: PendingActivities,
heartbeat_tracker: HeartbeatTracker,
grpc_liveness_waiters: crate::worker::GrpcLivenessWaiters,
drain_state: DrainState,
metrics: Option<Metrics>,
health: Option<HealthState>,
activity_mock_registry: Option<ActivityMockRegistry>,
outbox_store: Option<Arc<dyn OutboxStore>>,
namespace_store: Arc<dyn NamespaceStore>,
worker_deployment_store: Arc<dyn WorkerDeploymentStore>,
worker_supervisor: Arc<WorkerSupervisor>,
outbox_wake: Arc<tokio::sync::Notify>,
cluster_publisher: crate::cluster_publisher::ClusterEventPublisher,
transcript_publisher: crate::activity_publisher::ActivityEventPublisher,
attempt_owners: crate::worker::AttemptOwnerIndex,
queue_service_state: crate::worker::QueueServiceState,
queue_declarations: crate::worker::QueueDeclarationSource,
update_status: crate::update_check::UpdateStatusState,
workspace_root: crate::worker::WorkspaceRoot,
cluster_self_node: Option<String>,
cluster_responder: Option<aion_store_haematite::ClusterResponder>,
cluster_store: Option<Arc<aion_store_haematite::HaematiteStore>>,
watched_peers: Vec<crate::cluster::WatchedPeer>,
shard_directory: Option<Arc<crate::routing::StaticShardDirectory>>,
request_forwarder: Option<Arc<dyn crate::routing::RequestForwarder>>,
#[cfg(feature = "auth")]
jwks_cache: Option<JwksCache>,
}
impl ServerState {
const FALLBACK_CLUSTER_BROADCAST_CAPACITY: std::num::NonZeroUsize =
match std::num::NonZeroUsize::new(64) {
Some(value) => value,
None => std::num::NonZeroUsize::MIN,
};
pub async fn build(config: ServerConfig) -> Result<Self, ServerError> {
let (store_config, runtime) = config.into_parts();
let connected = connect_store(store_config).await?;
Self::build_with_connected_store(connected, runtime).await
}
pub async fn build_with_store<S>(store: S, runtime: RuntimeConfig) -> Result<Self, ServerError>
where
S: EventStore + NamespaceStore + WorkerDeploymentStore,
{
let leaf = Arc::new(store);
let namespace_store: Arc<dyn NamespaceStore> = leaf.clone();
let worker_deployment_store: Arc<dyn WorkerDeploymentStore> = leaf.clone();
Self::build_with_connected_store(
ConnectedStore::local(leaf, None, namespace_store, worker_deployment_store),
runtime,
)
.await
}
async fn build_with_connected_store(
connected: ConnectedStore,
runtime: RuntimeConfig,
) -> Result<Self, ServerError> {
let cluster_self_node = connected.cluster_self_node();
let outbox_store = connected.outbox_store;
let bootstrap_coordinator = connected.bootstrap_coordinator;
let cluster_responder = connected.cluster_responder;
let cluster_store = connected.cluster_store;
let watched_peers = connected.watched_peers;
let RoutingState {
shard_directory,
request_forwarder,
mint_routing,
} = build_routing_state(
cluster_store.as_ref(),
connected.directory_peers,
connected.self_node_id,
);
let (event_broadcast_capacity, query_timeout) = required_engine_seams(&runtime)?;
let (cluster_publisher, transcript_publisher) =
build_real_time_publishers(&runtime, connected.observability_store)?;
let (metrics, outbox_wake, instrumented_store) =
build_instrumented_store(&runtime, connected.event_store)?;
let exported_metrics = runtime.metrics.enabled.then_some(metrics.clone());
let seams = build_worker_seams(
&runtime,
&cluster_publisher,
&connected.namespace_store,
&connected.worker_deployment_store,
mint_routing,
);
let (
activity_dispatcher,
activity_mock_registry,
attempt_owners,
workspace_root,
update_status,
) = build_decorated_dispatcher(&runtime, &seams, transcript_publisher.clone());
let engine = boot_engine(EngineAssembly {
seams: &seams,
instrumented_store: &instrumented_store,
event_broadcast_capacity,
query_timeout,
activity_dispatcher,
active_registry: Arc::new(aion::Registry::default()),
bootstrap_coordinator,
runtime: &runtime,
})
.await?;
let resolver = NamespaceResolver::from_config(runtime.namespace.clone(), engine);
let worker_supervisor =
new_supervisor(&connected.worker_deployment_store, &cluster_publisher);
#[cfg(feature = "auth")]
let jwks_cache = build_jwks_cache(&runtime).await?;
Ok(Self {
inner: Arc::new(ServerStateInner {
namespace_guard: NamespaceGuard::new(resolver),
runtime,
metrics: exported_metrics,
worker_registry: seams.worker_registry,
pending_activities: seams.pending_activities,
heartbeat_tracker: seams.heartbeat_tracker,
grpc_liveness_waiters: crate::worker::GrpcLivenessWaiters::new(),
drain_state: seams.drain_state,
health: Some(HealthState::new(instrumented_store, true)),
activity_mock_registry,
outbox_store,
namespace_store: connected.namespace_store,
worker_supervisor,
worker_deployment_store: connected.worker_deployment_store,
outbox_wake,
cluster_publisher,
transcript_publisher,
attempt_owners,
queue_service_state: seams.queue_service_state,
queue_declarations: seams.queue_declarations,
update_status,
workspace_root,
cluster_self_node,
cluster_responder,
cluster_store,
watched_peers,
shard_directory,
request_forwarder,
#[cfg(feature = "auth")]
jwks_cache,
}),
})
}
#[must_use]
pub fn from_parts(namespace_resolver: NamespaceResolver, runtime: RuntimeConfig) -> Self {
Self::from_parts_with_namespace_store(
namespace_resolver,
runtime,
Arc::new(aion_store::InMemoryStore::default()),
)
}
#[must_use]
pub fn from_parts_with_namespace_store<S>(
namespace_resolver: NamespaceResolver,
runtime: RuntimeConfig,
store: Arc<S>,
) -> Self
where
S: NamespaceStore + WorkerDeploymentStore,
{
let namespace_store: Arc<dyn NamespaceStore> = store.clone();
let worker_deployment_store: Arc<dyn WorkerDeploymentStore> = store;
Self::from_parts_with_control_stores(
namespace_resolver,
runtime,
namespace_store,
worker_deployment_store,
)
}
#[must_use]
pub fn from_parts_with_control_stores(
namespace_resolver: NamespaceResolver,
runtime: RuntimeConfig,
namespace_store: Arc<dyn NamespaceStore>,
worker_deployment_store: Arc<dyn WorkerDeploymentStore>,
) -> Self {
let heartbeat_tracker = HeartbeatTracker::new(runtime.worker.heartbeat_window);
let pending_activities =
PendingActivities::default().with_heartbeat_window(runtime.worker.heartbeat_window);
let bounds = transcript_bounds(&runtime);
let batch = required_transcript_batch_policy(&runtime)
.unwrap_or(crate::activity_publisher::TranscriptBatchPolicy::UNBATCHED);
let cluster_publisher = crate::cluster_publisher::ClusterEventPublisher::new(
Self::FALLBACK_CLUSTER_BROADCAST_CAPACITY,
);
Self {
inner: Arc::new(ServerStateInner {
namespace_guard: NamespaceGuard::new(namespace_resolver),
runtime,
worker_registry: ConnectedWorkerRegistry::default()
.with_worker_deployment_store(worker_deployment_store.clone())
.with_cluster_publisher(cluster_publisher.clone()),
pending_activities,
heartbeat_tracker,
grpc_liveness_waiters: crate::worker::GrpcLivenessWaiters::new(),
drain_state: DrainState::default(),
metrics: None,
health: None,
activity_mock_registry: None,
outbox_store: None,
namespace_store,
worker_supervisor: new_supervisor(&worker_deployment_store, &cluster_publisher),
worker_deployment_store,
outbox_wake: Arc::new(tokio::sync::Notify::new()),
cluster_publisher,
transcript_publisher: build_transcript_publisher(
None,
Self::FALLBACK_CLUSTER_BROADCAST_CAPACITY,
bounds,
batch,
),
attempt_owners: crate::worker::AttemptOwnerIndex::new(),
queue_service_state: crate::worker::QueueServiceState::default(),
queue_declarations: crate::worker::QueueDeclarationSource::default(),
update_status: crate::update_check::UpdateStatusState::default(),
workspace_root: crate::worker::WorkspaceRoot::resolve(),
cluster_self_node: None,
cluster_responder: None,
cluster_store: None,
watched_peers: Vec::new(),
shard_directory: None,
request_forwarder: None,
#[cfg(feature = "auth")]
jwks_cache: None,
}),
}
}
#[cfg(feature = "auth")]
#[must_use]
pub fn from_parts_with_namespace_store_and_jwks<S>(
namespace_resolver: NamespaceResolver,
runtime: RuntimeConfig,
store: Arc<S>,
jwks_cache: JwksCache,
) -> Self
where
S: NamespaceStore + WorkerDeploymentStore,
{
let namespace_store: Arc<dyn NamespaceStore> = store.clone();
let worker_deployment_store: Arc<dyn WorkerDeploymentStore> = store;
let heartbeat_tracker = HeartbeatTracker::new(runtime.worker.heartbeat_window);
let pending_activities =
PendingActivities::default().with_heartbeat_window(runtime.worker.heartbeat_window);
let bounds = transcript_bounds(&runtime);
let batch = required_transcript_batch_policy(&runtime)
.unwrap_or(crate::activity_publisher::TranscriptBatchPolicy::UNBATCHED);
let cluster_publisher = crate::cluster_publisher::ClusterEventPublisher::new(
Self::FALLBACK_CLUSTER_BROADCAST_CAPACITY,
);
Self {
inner: Arc::new(ServerStateInner {
namespace_guard: NamespaceGuard::new(namespace_resolver),
runtime,
worker_registry: ConnectedWorkerRegistry::default()
.with_worker_deployment_store(worker_deployment_store.clone())
.with_cluster_publisher(cluster_publisher.clone()),
pending_activities,
heartbeat_tracker,
grpc_liveness_waiters: crate::worker::GrpcLivenessWaiters::new(),
drain_state: DrainState::default(),
metrics: None,
health: None,
activity_mock_registry: None,
outbox_store: None,
namespace_store,
worker_supervisor: new_supervisor(&worker_deployment_store, &cluster_publisher),
worker_deployment_store,
outbox_wake: Arc::new(tokio::sync::Notify::new()),
cluster_publisher,
transcript_publisher: build_transcript_publisher(
None,
Self::FALLBACK_CLUSTER_BROADCAST_CAPACITY,
bounds,
batch,
),
attempt_owners: crate::worker::AttemptOwnerIndex::new(),
queue_service_state: crate::worker::QueueServiceState::default(),
queue_declarations: crate::worker::QueueDeclarationSource::default(),
update_status: crate::update_check::UpdateStatusState::default(),
workspace_root: crate::worker::WorkspaceRoot::resolve(),
cluster_self_node: None,
cluster_responder: None,
cluster_store: None,
watched_peers: Vec::new(),
shard_directory: None,
request_forwarder: None,
jwks_cache: Some(jwks_cache),
}),
}
}
#[cfg(feature = "auth")]
#[must_use]
pub fn from_parts_with_jwks(
namespace_resolver: NamespaceResolver,
runtime: RuntimeConfig,
jwks_cache: JwksCache,
) -> Self {
let fallback_deployment_store: Arc<dyn WorkerDeploymentStore> =
Arc::new(aion_store::InMemoryStore::default());
let heartbeat_tracker = HeartbeatTracker::new(runtime.worker.heartbeat_window);
let pending_activities =
PendingActivities::default().with_heartbeat_window(runtime.worker.heartbeat_window);
let bounds = transcript_bounds(&runtime);
let batch = required_transcript_batch_policy(&runtime)
.unwrap_or(crate::activity_publisher::TranscriptBatchPolicy::UNBATCHED);
let cluster_publisher = crate::cluster_publisher::ClusterEventPublisher::new(
Self::FALLBACK_CLUSTER_BROADCAST_CAPACITY,
);
Self {
inner: Arc::new(ServerStateInner {
namespace_guard: NamespaceGuard::new(namespace_resolver),
runtime,
worker_registry: ConnectedWorkerRegistry::default(),
pending_activities,
heartbeat_tracker,
grpc_liveness_waiters: crate::worker::GrpcLivenessWaiters::new(),
drain_state: DrainState::default(),
metrics: None,
health: None,
activity_mock_registry: None,
outbox_store: None,
namespace_store: Arc::new(aion_store::InMemoryStore::default()),
worker_supervisor: new_supervisor(&fallback_deployment_store, &cluster_publisher),
worker_deployment_store: fallback_deployment_store,
outbox_wake: Arc::new(tokio::sync::Notify::new()),
cluster_publisher,
transcript_publisher: build_transcript_publisher(
None,
Self::FALLBACK_CLUSTER_BROADCAST_CAPACITY,
bounds,
batch,
),
attempt_owners: crate::worker::AttemptOwnerIndex::new(),
queue_service_state: crate::worker::QueueServiceState::default(),
queue_declarations: crate::worker::QueueDeclarationSource::default(),
update_status: crate::update_check::UpdateStatusState::default(),
workspace_root: crate::worker::WorkspaceRoot::resolve(),
cluster_self_node: None,
cluster_responder: None,
cluster_store: None,
watched_peers: Vec::new(),
shard_directory: None,
request_forwarder: None,
jwks_cache: Some(jwks_cache),
}),
}
}
#[must_use]
pub fn from_parts_with_registry(
namespace_resolver: NamespaceResolver,
runtime: RuntimeConfig,
worker_registry: ConnectedWorkerRegistry,
) -> Self {
let fallback_deployment_store: Arc<dyn WorkerDeploymentStore> =
Arc::new(aion_store::InMemoryStore::default());
let heartbeat_tracker = HeartbeatTracker::new(runtime.worker.heartbeat_window);
let pending_activities =
PendingActivities::default().with_heartbeat_window(runtime.worker.heartbeat_window);
let bounds = transcript_bounds(&runtime);
let batch = required_transcript_batch_policy(&runtime)
.unwrap_or(crate::activity_publisher::TranscriptBatchPolicy::UNBATCHED);
let cluster_publisher = crate::cluster_publisher::ClusterEventPublisher::new(
Self::FALLBACK_CLUSTER_BROADCAST_CAPACITY,
);
Self {
inner: Arc::new(ServerStateInner {
namespace_guard: NamespaceGuard::new(namespace_resolver),
runtime,
worker_registry,
pending_activities,
heartbeat_tracker,
grpc_liveness_waiters: crate::worker::GrpcLivenessWaiters::new(),
drain_state: DrainState::default(),
metrics: None,
health: None,
activity_mock_registry: None,
outbox_store: None,
namespace_store: Arc::new(aion_store::InMemoryStore::default()),
worker_supervisor: new_supervisor(&fallback_deployment_store, &cluster_publisher),
worker_deployment_store: fallback_deployment_store,
outbox_wake: Arc::new(tokio::sync::Notify::new()),
cluster_publisher,
transcript_publisher: build_transcript_publisher(
None,
Self::FALLBACK_CLUSTER_BROADCAST_CAPACITY,
bounds,
batch,
),
attempt_owners: crate::worker::AttemptOwnerIndex::new(),
queue_service_state: crate::worker::QueueServiceState::default(),
queue_declarations: crate::worker::QueueDeclarationSource::default(),
update_status: crate::update_check::UpdateStatusState::default(),
workspace_root: crate::worker::WorkspaceRoot::resolve(),
cluster_self_node: None,
cluster_responder: None,
cluster_store: None,
watched_peers: Vec::new(),
shard_directory: None,
request_forwarder: None,
#[cfg(feature = "auth")]
jwks_cache: None,
}),
}
}
#[must_use]
pub fn namespace_guard(&self) -> &NamespaceGuard {
&self.inner.namespace_guard
}
#[must_use]
pub fn deploy_guard(&self) -> crate::deploy::DeployGuard {
crate::deploy::DeployGuard::new(self.inner.namespace_guard.resolver().clone())
}
#[must_use]
pub fn runtime_config(&self) -> &RuntimeConfig {
&self.inner.runtime
}
#[must_use]
pub(crate) fn update_status(&self) -> &crate::update_check::UpdateStatusState {
&self.inner.update_status
}
#[must_use]
pub fn workspace_root(&self) -> &crate::worker::WorkspaceRoot {
&self.inner.workspace_root
}
#[must_use]
pub fn worker_registry(&self) -> &ConnectedWorkerRegistry {
&self.inner.worker_registry
}
pub fn cancel_in_flight_activities(
&self,
workflow_id: &aion_core::WorkflowId,
) -> Result<Vec<crate::worker::CancelRequest>, ServerError> {
crate::worker::cancel_in_flight_activities(
self.heartbeat_tracker(),
self.worker_registry(),
workflow_id,
)
}
#[must_use]
pub fn cluster_publisher(&self) -> &crate::cluster_publisher::ClusterEventPublisher {
&self.inner.cluster_publisher
}
#[must_use]
pub fn transcript_publisher(&self) -> &crate::activity_publisher::ActivityEventPublisher {
&self.inner.transcript_publisher
}
#[must_use]
pub fn attempt_owners(&self) -> &crate::worker::AttemptOwnerIndex {
&self.inner.attempt_owners
}
#[must_use]
pub fn queue_service_state(&self) -> &crate::worker::QueueServiceState {
&self.inner.queue_service_state
}
#[must_use]
pub fn queue_declarations(&self) -> &crate::worker::QueueDeclarationSource {
&self.inner.queue_declarations
}
pub fn unserved_queues(&self) -> Result<Vec<crate::worker::UnservedQueue>, ServerError> {
self.inner.queue_service_state.unserved()
}
pub fn unrecoverable_runs(
&self,
) -> Result<Vec<(aion_core::WorkflowId, aion::registry::UnrecoverableRun)>, ServerError> {
self.engine()?
.registry()
.unrecoverable()
.list()
.map_err(ServerError::from)
}
#[must_use]
pub fn intervention_router(&self) -> crate::worker::InterventionRouter {
let transport: std::sync::Arc<dyn crate::worker::InterventionTransport> = {
#[cfg(feature = "liminal-transport")]
{
std::sync::Arc::new(crate::worker::LiminalInterventionTransport)
}
#[cfg(not(feature = "liminal-transport"))]
{
std::sync::Arc::new(NullInterventionTransport)
}
};
crate::worker::InterventionRouter::new(
self.inner.worker_registry.clone(),
self.inner.attempt_owners.clone(),
transport,
)
.with_transcript_publisher(self.inner.transcript_publisher.clone())
}
#[must_use]
pub fn cluster_self_node(&self) -> Option<&str> {
self.inner.cluster_self_node.as_deref()
}
pub fn engine(&self) -> Result<Arc<aion::Engine>, ServerError> {
self.inner
.namespace_guard
.resolver()
.engine()
.map(Arc::clone)
}
#[must_use]
pub fn pending_activities(&self) -> &PendingActivities {
&self.inner.pending_activities
}
#[must_use]
pub fn heartbeat_tracker(&self) -> &HeartbeatTracker {
&self.inner.heartbeat_tracker
}
#[must_use]
pub fn grpc_liveness_waiters(&self) -> &crate::worker::GrpcLivenessWaiters {
&self.inner.grpc_liveness_waiters
}
#[must_use]
pub fn drain_state(&self) -> &DrainState {
&self.inner.drain_state
}
#[must_use]
pub fn metrics(&self) -> Option<&Metrics> {
self.inner.metrics.as_ref()
}
#[must_use]
pub fn health(&self) -> Option<&HealthState> {
self.inner.health.as_ref()
}
#[must_use]
pub fn activity_mock_registry(&self) -> Option<&ActivityMockRegistry> {
self.inner.activity_mock_registry.as_ref()
}
#[must_use]
pub fn outbox_store(&self) -> Option<Arc<dyn OutboxStore>> {
self.inner.outbox_store.clone()
}
#[must_use]
pub fn namespace_store(&self) -> &Arc<dyn NamespaceStore> {
&self.inner.namespace_store
}
#[must_use]
pub fn worker_deployment_store(&self) -> &Arc<dyn WorkerDeploymentStore> {
&self.inner.worker_deployment_store
}
#[must_use]
pub fn worker_supervisor(&self) -> &Arc<WorkerSupervisor> {
&self.inner.worker_supervisor
}
#[must_use]
pub fn namespace_minter(&self) -> NamespaceMinter {
let minter = NamespaceMinter::new(
Arc::clone(&self.inner.namespace_store),
self.inner.runtime.auto_create,
)
.with_cluster_publisher(self.inner.cluster_publisher.clone());
match self.namespace_routing() {
Some(routing) => minter.with_routing(routing),
None => minter,
}
}
#[must_use]
pub fn namespace_routing(&self) -> Option<crate::namespace::NamespaceRouting> {
build_namespace_routing(
self.cluster_store(),
self.shard_directory(),
self.request_forwarder(),
)
}
#[must_use]
pub fn outbox_wake(&self) -> Arc<tokio::sync::Notify> {
Arc::clone(&self.inner.outbox_wake)
}
#[must_use]
pub fn is_clustered(&self) -> bool {
self.inner.cluster_responder.is_some()
}
#[must_use]
pub fn cluster_store(&self) -> Option<&Arc<aion_store_haematite::HaematiteStore>> {
self.inner.cluster_store.as_ref()
}
#[must_use]
pub fn shard_directory(&self) -> Option<&Arc<crate::routing::StaticShardDirectory>> {
self.inner.shard_directory.as_ref()
}
#[must_use]
pub fn request_forwarder(&self) -> Option<&Arc<dyn crate::routing::RequestForwarder>> {
self.inner.request_forwarder.as_ref()
}
#[must_use]
pub fn spawn_heartbeat_sweeper(
&self,
shutdown: tokio::sync::watch::Receiver<bool>,
) -> tokio::task::JoinHandle<()> {
let sweeper = crate::worker::HeartbeatSweeper::new(
self.inner.heartbeat_tracker.clone(),
self.inner.worker_registry.clone(),
self.inner.pending_activities.clone(),
self.inner.drain_state.clone(),
self.inner.runtime.worker.heartbeat_window,
)
.with_queue_state(self.inner.queue_service_state.clone());
tokio::spawn(sweeper.run(shutdown))
}
#[cfg(feature = "liminal-transport")]
#[must_use]
pub fn spawn_liminal_liveness_probe(
&self,
notifier: std::sync::Arc<crate::worker::LiminalConnectionNotifier>,
shutdown: tokio::sync::watch::Receiver<bool>,
) -> tokio::task::JoinHandle<()> {
let probe = crate::worker::LivenessProbe::across_transports(
Some(notifier),
Some(self.inner.grpc_liveness_waiters.clone()),
self.inner.heartbeat_tracker.clone(),
self.inner.worker_registry.clone(),
self.inner.runtime.worker.heartbeat_window,
);
tokio::spawn(probe.run(shutdown))
}
pub fn spawn_cluster_supervisor(
&self,
config: crate::cluster::SupervisorConfig,
shutdown: tokio::sync::watch::Receiver<bool>,
) -> Result<bool, ServerError> {
let Some(cluster_store) = self.inner.cluster_store.clone() else {
return Ok(false);
};
if self.inner.watched_peers.is_empty() {
return Ok(false);
}
let engine = Arc::clone(self.inner.namespace_guard.resolver().engine()?);
let publisher = Arc::new(self.inner.cluster_publisher.clone());
let self_node = self.inner.cluster_self_node.clone().unwrap_or_default();
let adopter = Arc::new(crate::cluster::OutboxSettlingAdopter::new(
engine,
self.inner.outbox_store.clone(),
));
let supervisor = crate::cluster::ClusterSupervisor::new(
cluster_store,
adopter,
self.inner.watched_peers.clone(),
config,
)
.with_publisher(publisher, self_node);
if !supervisor.watches_any() {
return Ok(false);
}
tokio::spawn(supervisor.run(shutdown));
Ok(true)
}
#[cfg(feature = "auth")]
#[must_use]
pub fn jwks_cache(&self) -> Option<&JwksCache> {
self.inner.jwks_cache.as_ref()
}
pub fn shutdown(&self) -> Result<(), ServerError> {
self.inner.namespace_guard.resolver().shutdown_engine()
}
}
#[cfg(feature = "auth")]
async fn build_jwks_cache(runtime: &RuntimeConfig) -> Result<Option<JwksCache>, ServerError> {
if !runtime.auth.enabled {
return Ok(None);
}
let Some(url) = runtime.auth.jwks_url.clone() else {
return Err(ServerError::Config {
message: "auth.jwks_url must not be empty when auth.enabled is true".to_owned(),
});
};
let interval = std::time::Duration::from_secs(runtime.auth.jwks_refresh_seconds);
let cache = JwksCache::new(url, interval)
.await
.map_err(|error| ServerError::Config {
message: format!("auth jwks initial fetch failed: {error}"),
})?;
Ok(Some(cache))
}
fn metrics_config_error(error: &MetricsError) -> ServerError {
ServerError::Config {
message: error.to_string(),
}
}
struct EngineAssembly<'a> {
seams: &'a WorkerSeams,
instrumented_store: &'a Arc<InstrumentedEventStore>,
event_broadcast_capacity: std::num::NonZeroUsize,
query_timeout: std::time::Duration,
activity_dispatcher: Arc<dyn ActivityDispatcher>,
active_registry: Arc<aion::Registry>,
bootstrap_coordinator: bool,
runtime: &'a RuntimeConfig,
}
fn server_search_attribute_schema() -> Result<aion_core::SearchAttributeSchema, ServerError> {
let mut schema = aion_core::SearchAttributeSchema::new();
for (name, label) in [
(crate::namespace::NAMESPACE_ATTRIBUTE, "namespace"),
(crate::namespace::TASK_QUEUE_ATTRIBUTE, "task_queue"),
(crate::namespace::DISPLAY_NAME_ATTRIBUTE, "display_name"),
] {
schema
.register(name, aion_core::SearchAttributeType::String)
.map_err(|error| ServerError::Config {
message: format!("failed to register {label} search attribute: {error}"),
})?;
}
Ok(schema)
}
async fn boot_engine(assembly: EngineAssembly<'_>) -> Result<Arc<aion::Engine>, ServerError> {
let search_attribute_schema = server_search_attribute_schema()?;
let runtime = assembly.runtime;
let builder = EngineBuilder::new()
.store_arc(assembly.instrumented_store.clone())
.event_streaming(assembly.event_broadcast_capacity)
.in_memory_visibility()
.search_attribute_schema(search_attribute_schema)
.scheduler_threads(runtime.scheduler_threads)
.outbox_enabled(runtime.outbox.enabled)
.activity_dispatcher(assembly.activity_dispatcher)
.active_registry(assembly.active_registry)
.production_recovery_seam()
.defer_startup_recovery()
.signal_router_factory(|runtime: Arc<RuntimeHandle>, handoff| {
Arc::new(ConcreteSignalRouter::new(runtime, handoff)) as Arc<dyn SignalRouter>
})
.query_timeout(assembly.query_timeout)
.bootstrap_schedule_coordinator(assembly.bootstrap_coordinator)
.load_workflow_sources(runtime.workflow_packages.iter().map(PathBuf::as_path));
let builder = if runtime.owned_shards.is_empty() {
builder
} else {
builder.owned_shards(runtime.owned_shards.iter().copied())
};
let engine = Arc::new(builder.build().await.map_err(ServerError::from)?);
install_engine_backed_seams(assembly.seams, &engine, runtime.outbox.enabled);
engine
.run_startup_recovery()
.await
.map_err(ServerError::from)?;
Ok(engine)
}
fn required_engine_seams(
runtime: &RuntimeConfig,
) -> Result<(std::num::NonZeroUsize, std::time::Duration), ServerError> {
let event_broadcast_capacity = runtime
.websocket
.event_broadcast_capacity
.and_then(std::num::NonZeroUsize::new)
.ok_or_else(|| ServerError::Config {
message: crate::config::EVENT_BROADCAST_CAPACITY_REQUIRED.to_owned(),
})?;
let query_timeout = runtime
.query_timeout
.filter(|timeout| !timeout.is_zero())
.ok_or_else(|| ServerError::Config {
message: crate::config::QUERY_TIMEOUT_REQUIRED.to_owned(),
})?;
Ok((event_broadcast_capacity, query_timeout))
}
fn install_outbox_delivery(
pending_activities: &PendingActivities,
engine: &Arc<aion::Engine>,
outbox_enabled: bool,
) {
if outbox_enabled {
let callback = Arc::new(crate::worker::ServerOutboxDeliveryCallback::new(
Arc::clone(engine),
));
pending_activities.set_outbox_delivery(callback);
}
}
struct WorkerSeams {
worker_registry: ConnectedWorkerRegistry,
pending_activities: PendingActivities,
heartbeat_tracker: HeartbeatTracker,
drain_state: DrainState,
queue_declarations: crate::worker::QueueDeclarationSource,
queue_service_state: crate::worker::QueueServiceState,
declared_bodies: crate::worker::DeclaredBodySource,
cluster_publisher: crate::cluster_publisher::ClusterEventPublisher,
}
fn build_worker_seams(
runtime: &RuntimeConfig,
cluster_publisher: &crate::cluster_publisher::ClusterEventPublisher,
namespace_store: &Arc<dyn NamespaceStore>,
worker_deployment_store: &Arc<dyn WorkerDeploymentStore>,
mint_routing: Option<crate::namespace::NamespaceRouting>,
) -> WorkerSeams {
let worker_registry = ConnectedWorkerRegistry::default()
.with_cluster_publisher(cluster_publisher.clone())
.with_namespace_minting(namespace_store.clone(), runtime.auto_create)
.with_worker_deployment_store(worker_deployment_store.clone());
let worker_registry = match mint_routing {
Some(routing) => worker_registry.with_namespace_routing(routing),
None => worker_registry,
};
WorkerSeams {
worker_registry,
pending_activities: PendingActivities::default()
.with_heartbeat_window(runtime.worker.heartbeat_window),
heartbeat_tracker: HeartbeatTracker::new(runtime.worker.heartbeat_window),
drain_state: DrainState::default(),
queue_declarations: crate::worker::QueueDeclarationSource::default(),
queue_service_state: crate::worker::QueueServiceState::default(),
declared_bodies: crate::worker::DeclaredBodySource::default(),
cluster_publisher: cluster_publisher.clone(),
}
}
fn install_queue_declarations(
queue_declarations: &crate::worker::QueueDeclarationSource,
engine: &Arc<aion::Engine>,
) {
queue_declarations.install(Arc::new(crate::worker::EngineQueueDeclarations::new(
Arc::clone(engine),
)));
}
fn install_declared_bodies(
declared_bodies: &crate::worker::DeclaredBodySource,
engine: &Arc<aion::Engine>,
) {
declared_bodies.install(Arc::new(crate::worker::EngineDeclaredBodies::new(
Arc::clone(engine),
)));
}
fn install_engine_backed_seams(seams: &WorkerSeams, engine: &Arc<aion::Engine>, outbox: bool) {
install_outbox_delivery(&seams.pending_activities, engine, outbox);
install_queue_declarations(&seams.queue_declarations, engine);
install_declared_bodies(&seams.declared_bodies, engine);
}
fn build_bridge_dispatcher(
runtime: &RuntimeConfig,
seams: &WorkerSeams,
) -> (WorkerActivityDispatcher, crate::worker::AttemptOwnerIndex) {
let attempt_owners = crate::worker::AttemptOwnerIndex::new();
let dispatcher = WorkerActivityDispatcher::new(
seams.worker_registry.clone(),
runtime.default_namespace.clone(),
seams.heartbeat_tracker.clone(),
)
.with_pending(seams.pending_activities.clone())
.with_drain_state(seams.drain_state.clone())
.with_tokio_handle(tokio::runtime::Handle::current())
.with_attempt_owners(attempt_owners.clone())
.with_queue_service(runtime.worker.queue_service.clone())
.with_queue_declarations(seams.queue_declarations.clone())
.with_queue_state(seams.queue_service_state.clone())
.with_cluster_publisher(seams.cluster_publisher.clone());
(dispatcher, attempt_owners)
}
fn build_instrumented_store(
runtime: &RuntimeConfig,
event_store: Arc<dyn EventStore>,
) -> Result<
(
Metrics,
Arc<tokio::sync::Notify>,
Arc<InstrumentedEventStore>,
),
ServerError,
> {
let metrics = Metrics::new().map_err(|error| metrics_config_error(&error))?;
let outbox_wake = Arc::new(tokio::sync::Notify::new());
let instrumented_store = Arc::new(
InstrumentedEventStore::new(
event_store,
metrics.clone(),
runtime.default_namespace.clone(),
)
.with_outbox_wake(Arc::clone(&outbox_wake)),
);
Ok((metrics, outbox_wake, instrumented_store))
}
fn build_decorated_dispatcher(
runtime: &RuntimeConfig,
seams: &WorkerSeams,
transcript: crate::activity_publisher::ActivityEventPublisher,
) -> (
Arc<dyn ActivityDispatcher>,
Option<ActivityMockRegistry>,
crate::worker::AttemptOwnerIndex,
crate::worker::WorkspaceRoot,
crate::update_check::UpdateStatusState,
) {
let workspace_root = crate::worker::WorkspaceRoot::resolve();
let update_status = crate::update_check::UpdateStatusState::default();
let (dispatcher, attempt_owners) = build_bridge_dispatcher(runtime, seams);
let (activity_dispatcher, activity_mock_registry) = decorate_activity_dispatcher(
dispatcher,
seams.declared_bodies.clone(),
workspace_root.clone(),
transcript,
update_status.clone(),
runtime.dev.enabled,
);
(
activity_dispatcher,
activity_mock_registry,
attempt_owners,
workspace_root,
update_status,
)
}
fn decorate_activity_dispatcher(
dispatcher: WorkerActivityDispatcher,
declared_bodies: crate::worker::DeclaredBodySource,
workspace_root: crate::worker::WorkspaceRoot,
transcript: crate::activity_publisher::ActivityEventPublisher,
update_status: crate::update_check::UpdateStatusState,
dev_enabled: bool,
) -> (Arc<dyn ActivityDispatcher>, Option<ActivityMockRegistry>) {
let declared = crate::worker::DeclaredCommandDispatcher::new(
Arc::new(dispatcher),
declared_bodies.clone(),
tokio::runtime::Handle::current(),
workspace_root,
transcript,
);
let observed = crate::update_check::UpdateCheckObserver::new(
Arc::new(declared),
declared_bodies,
update_status,
);
if dev_enabled {
let registry = ActivityMockRegistry::new();
let decorated = DevMockingDispatcher::new(Arc::new(observed), registry.clone());
(Arc::new(decorated), Some(registry))
} else {
(Arc::new(observed), None)
}
}
fn required_cluster_broadcast_capacity(
runtime: &RuntimeConfig,
) -> Result<std::num::NonZeroUsize, ServerError> {
runtime
.websocket
.cluster_broadcast_capacity
.and_then(std::num::NonZeroUsize::new)
.ok_or_else(|| ServerError::Config {
message: crate::config::CLUSTER_BROADCAST_CAPACITY_REQUIRED.to_owned(),
})
}
fn build_real_time_publishers(
runtime: &RuntimeConfig,
observability_store: Option<Arc<dyn aion_store::ObservabilityStore>>,
) -> Result<
(
crate::cluster_publisher::ClusterEventPublisher,
crate::activity_publisher::ActivityEventPublisher,
),
ServerError,
> {
let capacity = required_cluster_broadcast_capacity(runtime)?;
let batch = required_transcript_batch_policy(runtime)?;
Ok((
crate::cluster_publisher::ClusterEventPublisher::new(capacity),
build_transcript_publisher(
observability_store,
capacity,
transcript_bounds(runtime),
batch,
),
))
}
fn required_node_cache_budget(
config: &StoreConfig,
) -> Result<haematite::NodeCacheBudget, ServerError> {
config.node_cache_budget.ok_or_else(|| ServerError::Config {
message: crate::config::STORE_NODE_CACHE_BUDGET_REQUIRED.to_owned(),
})
}
fn required_transcript_batch_policy(
runtime: &RuntimeConfig,
) -> Result<crate::activity_publisher::TranscriptBatchPolicy, ServerError> {
let max_batch_events = runtime
.observability
.max_batch_events
.and_then(std::num::NonZeroUsize::new)
.ok_or_else(|| ServerError::Config {
message: crate::config::OBSERVABILITY_MAX_BATCH_EVENTS_REQUIRED.to_owned(),
})?;
let max_batch_hold_ms =
runtime
.observability
.max_batch_hold_ms
.ok_or_else(|| ServerError::Config {
message: crate::config::OBSERVABILITY_MAX_BATCH_HOLD_MS_REQUIRED.to_owned(),
})?;
Ok(crate::activity_publisher::TranscriptBatchPolicy {
max_batch_events,
max_hold: std::time::Duration::from_millis(max_batch_hold_ms),
})
}
fn transcript_bounds(runtime: &RuntimeConfig) -> crate::activity_bounds::TranscriptBounds {
crate::activity_bounds::TranscriptBounds {
max_event_bytes: runtime.observability.max_event_bytes,
max_stream_events: runtime.observability.max_stream_events,
}
}
fn build_transcript_publisher(
observability_store: Option<Arc<dyn aion_store::ObservabilityStore>>,
capacity: std::num::NonZeroUsize,
bounds: crate::activity_bounds::TranscriptBounds,
batch: crate::activity_publisher::TranscriptBatchPolicy,
) -> crate::activity_publisher::ActivityEventPublisher {
let store = observability_store
.unwrap_or_else(|| Arc::new(aion_store::InMemoryObservabilityStore::default()));
crate::activity_publisher::ActivityEventPublisher::new(store, capacity, batch)
.with_bounds(bounds)
}
struct RoutingState {
shard_directory: Option<Arc<crate::routing::StaticShardDirectory>>,
request_forwarder: Option<Arc<dyn crate::routing::RequestForwarder>>,
mint_routing: Option<crate::namespace::NamespaceRouting>,
}
fn build_routing_state(
cluster_store: Option<&Arc<aion_store_haematite::HaematiteStore>>,
directory_peers: Vec<crate::routing::DirectoryPeer>,
self_node_id: Option<String>,
) -> RoutingState {
let Some(store) = cluster_store else {
return RoutingState {
shard_directory: None,
request_forwarder: None,
mint_routing: None,
};
};
let shard_directory = Arc::new(crate::routing::StaticShardDirectory::new(
Arc::clone(store),
directory_peers,
self_node_id,
));
let request_forwarder: Arc<dyn crate::routing::RequestForwarder> =
Arc::new(crate::routing::GrpcRequestForwarder::new());
let mint_routing = build_namespace_routing(
Some(store),
Some(&shard_directory),
Some(&request_forwarder),
);
RoutingState {
shard_directory: Some(shard_directory),
request_forwarder: Some(request_forwarder),
mint_routing,
}
}
fn build_namespace_routing(
cluster_store: Option<&Arc<aion_store_haematite::HaematiteStore>>,
shard_directory: Option<&Arc<crate::routing::StaticShardDirectory>>,
request_forwarder: Option<&Arc<dyn crate::routing::RequestForwarder>>,
) -> Option<crate::namespace::NamespaceRouting> {
use crate::namespace::{
GrpcMintForwarder, MintForwarder, MintShardOwners, NamespaceRouting, NamespaceShardResolver,
};
let store = Arc::clone(cluster_store?);
let directory = Arc::clone(shard_directory?);
let shards: Arc<dyn NamespaceShardResolver> = store;
let owners: Arc<dyn MintShardOwners> = directory;
let forwarder: Arc<dyn MintForwarder> =
Arc::new(GrpcMintForwarder::new(Arc::clone(request_forwarder?)));
Some(NamespaceRouting::new(shards, owners, forwarder))
}
struct ConnectedStore {
event_store: Arc<dyn EventStore>,
outbox_store: Option<Arc<dyn OutboxStore>>,
namespace_store: Arc<dyn NamespaceStore>,
worker_deployment_store: Arc<dyn WorkerDeploymentStore>,
observability_store: Option<Arc<dyn aion_store::ObservabilityStore>>,
bootstrap_coordinator: bool,
cluster_responder: Option<aion_store_haematite::ClusterResponder>,
cluster_store: Option<Arc<aion_store_haematite::HaematiteStore>>,
watched_peers: Vec<crate::cluster::WatchedPeer>,
directory_peers: Vec<crate::routing::DirectoryPeer>,
self_node_id: Option<String>,
}
impl ConnectedStore {
fn local(
event_store: Arc<dyn EventStore>,
outbox_store: Option<Arc<dyn OutboxStore>>,
namespace_store: Arc<dyn NamespaceStore>,
worker_deployment_store: Arc<dyn WorkerDeploymentStore>,
) -> Self {
Self {
event_store,
outbox_store,
namespace_store,
worker_deployment_store,
observability_store: None,
bootstrap_coordinator: true,
cluster_responder: None,
cluster_store: None,
watched_peers: Vec::new(),
directory_peers: Vec::new(),
self_node_id: None,
}
}
fn cluster_self_node(&self) -> Option<String> {
self.self_node_id.clone()
}
}
async fn connect_store(config: StoreConfig) -> Result<ConnectedStore, ServerError> {
match config.backend {
StoreBackend::Memory => {
let leaf = Arc::new(aion_store::InMemoryStore::default());
let namespace_store: Arc<dyn NamespaceStore> = leaf.clone();
let worker_deployment_store: Arc<dyn WorkerDeploymentStore> = leaf.clone();
Ok(ConnectedStore::local(
leaf,
None,
namespace_store,
worker_deployment_store,
))
}
StoreBackend::Haematite => connect_haematite_store(config).await,
}
}
async fn connect_haematite_store(config: StoreConfig) -> Result<ConnectedStore, ServerError> {
let node_cache_budget = required_node_cache_budget(&config)?;
let Some(data_dir) = config.data_dir else {
return Err(ServerError::Config {
message: "store.data_dir must not be empty when store.backend is haematite".to_owned(),
});
};
let shard_count = config.shard_count;
let owned_shards = config.owned_shards.clone();
let cluster = config.cluster.clone();
let watched_peers: Vec<crate::cluster::WatchedPeer> = cluster
.as_ref()
.map(|cluster| {
cluster
.peers
.iter()
.map(|peer| crate::cluster::WatchedPeer {
name: peer.name.clone(),
owned_shards: peer.owned_shards.clone(),
})
.collect()
})
.unwrap_or_default();
let directory_peers: Vec<crate::routing::DirectoryPeer> = cluster
.as_ref()
.map(|cluster| {
cluster
.peers
.iter()
.map(|peer| crate::routing::DirectoryPeer {
name: peer.name.clone(),
owned_shards: peer.owned_shards.clone(),
grpc_addr: peer.grpc_address,
})
.collect()
})
.unwrap_or_default();
let self_node_id: Option<String> = cluster.as_ref().map(|cluster| cluster.node_id.clone());
let (store, responder) = tokio::task::spawn_blocking(move || {
build_haematite_store(&data_dir, shard_count, cluster, node_cache_budget)
})
.await
.map_err(|error| ServerError::Config {
message: format!("haematite store initialization task failed: {error}"),
})??;
let bootstrap_coordinator = if owned_shards.is_empty() {
true
} else {
store.set_owned_shards(owned_shards.iter().copied());
store.owns_workflow_shard(&aion::schedule_coordinator_workflow_id())
};
let leaf = Arc::new(store);
let event_store: Arc<dyn EventStore> = leaf.clone();
let outbox_store: Arc<dyn OutboxStore> = leaf.clone();
let namespace_store: Arc<dyn NamespaceStore> = leaf.clone();
let worker_deployment_store: Arc<dyn WorkerDeploymentStore> = leaf.clone();
let observability_store: Arc<dyn aion_store::ObservabilityStore> = leaf.clone();
let cluster_store = responder.as_ref().map(|_| leaf);
let (watched_peers, directory_peers, self_node_id) = if cluster_store.is_some() {
(watched_peers, directory_peers, self_node_id)
} else {
(Vec::new(), Vec::new(), None)
};
Ok(ConnectedStore {
event_store,
outbox_store: Some(outbox_store),
namespace_store,
worker_deployment_store,
observability_store: Some(observability_store),
bootstrap_coordinator,
cluster_responder: responder,
cluster_store,
watched_peers,
directory_peers,
self_node_id,
})
}
fn build_haematite_store(
data_dir: &str,
shard_count: usize,
cluster: Option<crate::config::ClusterConfig>,
node_cache_budget: haematite::NodeCacheBudget,
) -> Result<
(
aion_store_haematite::HaematiteStore,
Option<aion_store_haematite::ClusterResponder>,
),
ServerError,
> {
build_haematite_store_with_hook(data_dir, shard_count, cluster, node_cache_budget, || Ok(()))
}
fn build_haematite_store_with_hook(
data_dir: &str,
shard_count: usize,
cluster: Option<crate::config::ClusterConfig>,
node_cache_budget: haematite::NodeCacheBudget,
before_backend_touch: impl FnOnce() -> Result<(), std::io::Error>,
) -> Result<
(
aion_store_haematite::HaematiteStore,
Option<aion_store_haematite::ClusterResponder>,
),
ServerError,
> {
use aion_store_haematite::{ClusterBootstrap, HaematiteStore};
let private_root = crate::filesystem::ConfinedDir::open_or_create(std::path::Path::new(
data_dir,
))
.map_err(|error| ServerError::Config {
message: format!("unsafe store.data_dir `{data_dir}`: {error}"),
})?;
for shard in 0..shard_count {
private_root
.create_dir_all(std::path::Path::new(&format!("shard-{shard}")))
.map_err(|error| ServerError::Config {
message: format!(
"failed to materialize shard-{shard} under store.data_dir `{data_dir}`: {error}"
),
})?;
}
private_root
.harden_tree()
.map_err(|error| private_store_mode_error(data_dir, &error))?;
before_backend_touch().map_err(|error| ServerError::Config {
message: format!("store.data_dir pre-open hook failed: {error}"),
})?;
#[cfg(unix)]
let backend_path = private_root
.backend_path()
.map_err(|error| ServerError::Config {
message: format!("failed to resolve held store.data_dir `{data_dir}`: {error}"),
})?;
#[cfg(all(unix, not(any(target_os = "linux", target_os = "android"))))]
crate::filesystem::validate_ambient_backend_ancestors(&backend_path).map_err(|error| {
let (component, reason) = error.into_parts();
ServerError::UnsafeDataRootAncestor {
data_root: backend_path.clone(),
component,
reason,
}
})?;
#[cfg(not(unix))]
let backend_path = std::path::PathBuf::from(data_dir);
let Some(cluster) = cluster else {
let store = if backend_path.join("config.json").exists() {
HaematiteStore::open(&backend_path, node_cache_budget).map_err(ServerError::from)?
} else {
HaematiteStore::create_with_shard_count(&backend_path, shard_count, node_cache_budget)
.map_err(ServerError::from)?
};
store.materialize_all_shards().map_err(ServerError::from)?;
private_root
.harden_tree()
.map_err(|error| private_store_mode_error(data_dir, &error))?;
let store = store.retain_data_root_capability(private_root);
return Ok((store, None));
};
let boot = ClusterBootstrap {
node_id: cluster.node_id,
bind_address: cluster.bind_address,
members: cluster.members,
peers: cluster
.peers
.into_iter()
.map(|peer| (peer.name, peer.address))
.collect(),
timeout: HAEMATITE_CLUSTER_OP_TIMEOUT,
};
let (store, responder) = HaematiteStore::open_or_create_distributed(
&backend_path,
shard_count,
boot,
node_cache_budget,
)
.map_err(ServerError::from)?;
store.materialize_all_shards().map_err(ServerError::from)?;
private_root
.harden_tree()
.map_err(|error| private_store_mode_error(data_dir, &error))?;
let store = store.retain_data_root_capability(private_root);
Ok((store, Some(responder)))
}
fn private_store_mode_error(data_dir: &str, error: &std::io::Error) -> ServerError {
ServerError::Config {
message: format!(
"failed to apply private modes under store.data_dir `{data_dir}`: {error}"
),
}
}
const HAEMATITE_CLUSTER_OP_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(5);
#[cfg(not(feature = "liminal-transport"))]
#[derive(Clone, Debug)]
struct NullInterventionTransport;
#[cfg(not(feature = "liminal-transport"))]
#[async_trait::async_trait]
impl crate::worker::InterventionTransport for NullInterventionTransport {
async fn push(
&self,
_worker: &crate::worker::WorkerHandle,
_command: aion_core::InterventionCommand,
) -> Result<aion_core::InterventionOutcome, ServerError> {
Err(ServerError::worker_connection_lost(
"intervention",
"no intervention push transport is compiled in".to_owned(),
))
}
}
#[cfg(test)]
mod tests {
use std::{net::SocketAddr, time::Duration};
use aion_store::InMemoryStore;
use super::ServerState;
use crate::config::{
AuthConfig, AuthoringConfig, DeployConfig, DevConfig, ListenConfig, MetricsConfig,
NamespaceConfig, NamespaceMode, OpsConsoleAssetSource, OpsConsoleConfig, OutboxConfig,
RuntimeConfig, WebSocketConfig, WorkerConfig,
};
fn runtime_config() -> RuntimeConfig {
RuntimeConfig {
listen: ListenConfig {
grpc: SocketAddr::from(([127, 0, 0, 1], 50051)),
http: SocketAddr::from(([127, 0, 0, 1], 8080)),
},
tls: None,
auth: AuthConfig {
enabled: false,
jwks_url: None,
jwks_refresh_seconds: 300,
},
ops_console: OpsConsoleConfig {
source: OpsConsoleAssetSource::Embedded,
},
namespace: NamespaceConfig {
mode: NamespaceMode::SharedEngine,
},
worker: WorkerConfig {
heartbeat_window: Duration::from_secs(30),
..WorkerConfig::default()
},
websocket: WebSocketConfig {
outbound_buffer_bound: 32,
event_broadcast_capacity: Some(64),
cluster_broadcast_capacity: Some(64),
},
workflow_packages: Vec::new(),
deploy: DeployConfig::default(),
authoring: AuthoringConfig::default(),
dev: DevConfig::default(),
outbox: OutboxConfig::default(),
observability: crate::config::ObservabilityConfig::with_flush_policy(64, 0),
mcp: crate::config::ResolvedMcpConfig::default(),
scheduler_threads: 1,
query_timeout: Some(Duration::from_secs(10)),
default_namespace: "default".to_owned(),
auto_create: crate::config::AutoCreate::Open,
max_in_flight_activities: crate::config::DEFAULT_MAX_IN_FLIGHT_ACTIVITIES,
drain_timeout: Duration::from_secs(30),
metrics: MetricsConfig { enabled: true },
owned_shards: Vec::new(),
cors_allowed_origins: Vec::new(),
}
}
#[test]
fn an_unruled_flush_policy_refuses_to_build_the_publisher() {
for (mutate, expected_key) in [
(
Box::new(|runtime: &mut RuntimeConfig| {
runtime.observability.max_batch_events = None;
}) as Box<dyn Fn(&mut RuntimeConfig)>,
"observability.max_batch_events",
),
(
Box::new(|runtime: &mut RuntimeConfig| {
runtime.observability.max_batch_events = Some(0);
}),
"observability.max_batch_events",
),
(
Box::new(|runtime: &mut RuntimeConfig| {
runtime.observability.max_batch_hold_ms = None;
}),
"observability.max_batch_hold_ms",
),
] {
let mut runtime = runtime_config();
mutate(&mut runtime);
let message = super::required_transcript_batch_policy(&runtime)
.err()
.map_or_else(String::new, |error| error.to_string());
assert!(
message.contains(expected_key),
"the refusal must name {expected_key}: {message}"
);
assert!(
message.contains("no default"),
"and must say the key has no default: {message}"
);
}
}
#[test]
fn a_stated_flush_policy_is_carried_through_verbatim() -> Result<(), Box<dyn std::error::Error>>
{
let mut runtime = runtime_config();
runtime.observability = crate::config::ObservabilityConfig::with_flush_policy(32, 0);
let policy = super::required_transcript_batch_policy(&runtime)?;
assert_eq!(policy.max_batch_events.get(), 32);
assert_eq!(policy.max_hold, Duration::ZERO);
runtime.observability = crate::config::ObservabilityConfig::with_flush_policy(8, 250);
let policy = super::required_transcript_batch_policy(&runtime)?;
assert_eq!(policy.max_batch_events.get(), 8);
assert_eq!(policy.max_hold, Duration::from_millis(250));
Ok(())
}
#[test]
fn engine_schema_accepts_every_attribute_the_start_writer_records()
-> Result<(), Box<dyn std::error::Error>> {
let schema = super::server_search_attribute_schema()?;
let recorded = crate::api::handlers::workflows::start_search_attributes(
"tenant-a",
Some("gpu"),
Some("Nightly settlement"),
);
assert!(
recorded.contains_key(crate::namespace::DISPLAY_NAME_ATTRIBUTE),
"the fixture must exercise the display-name attribute, or this test \
cannot see its registration go missing"
);
for (name, value) in &recorded {
schema.validate(name, value).map_err(|error| {
format!(
"the start writer records {name}, but the engine's schema refuses it: {error}"
)
})?;
}
Ok(())
}
#[tokio::test]
async fn builds_state_with_in_memory_store() -> Result<(), Box<dyn std::error::Error>> {
let state =
ServerState::build_with_store(InMemoryStore::default(), runtime_config()).await?;
std::hint::black_box(state.namespace_guard());
std::hint::black_box(state.worker_registry());
Ok(())
}
#[tokio::test]
async fn unserved_queues_surfaces_a_parked_dispatch_and_clears_it()
-> Result<(), Box<dyn std::error::Error>> {
use aion::{ActivityDispatch, ActivityDispatcher as _};
use aion_core::{ActivityId, RunId, WorkflowId};
use std::sync::Arc;
let state =
ServerState::build_with_store(InMemoryStore::default(), runtime_config()).await?;
assert!(
state.unserved_queues()?.is_empty(),
"a calm boot has no unserved queues"
);
assert!(state.queue_declarations().is_installed());
assert_eq!(
state
.queue_declarations()
.declaration_for("nobody-serves-this"),
crate::worker::QueueDeclaration::Unknown
);
let dispatcher = Arc::new(
crate::worker::WorkerActivityDispatcher::new(
state.worker_registry().clone(),
"default",
crate::worker::HeartbeatTracker::new(Duration::from_secs(5)),
)
.with_queue_state(state.queue_service_state().clone())
.with_queue_declarations(state.queue_declarations().clone()),
);
let workflow_id = WorkflowId::new_v4();
let request = ActivityDispatch {
namespace: "default".to_owned(),
task_queue: "nobody-serves-this".to_owned(),
node: None,
workflow_id: workflow_id.clone(),
run_id: RunId::new_v4(),
activity_id: ActivityId::from_sequence_position(0),
name: "greet".to_owned(),
input: "{}".to_owned(),
config: "{}".to_owned(),
attempt: 1,
advisory: false,
labels: std::collections::BTreeMap::new(),
};
let parked = std::thread::spawn(move || dispatcher.dispatch(request));
let mut unserved = Vec::new();
for _ in 0..30 {
unserved = state.unserved_queues()?;
if !unserved.is_empty() {
break;
}
tokio::time::sleep(Duration::from_millis(100)).await;
}
assert_eq!(unserved.len(), 1, "the parked dispatch is not surfaced");
assert_eq!(
unserved[0].reason,
crate::worker::QueueServiceReason::NoLivePollers,
"an empty catalog must not be read as a structural refusal"
);
assert_eq!(unserved[0].key.task_queue, "nobody-serves-this");
assert_eq!(unserved[0].waiting.len(), 1);
assert_eq!(unserved[0].waiting[0].workflow_id, workflow_id);
let (worker_tx, worker_rx) = tokio::sync::mpsc::channel(1);
drop(worker_rx);
let registration = state.worker_registry().register_namespaces(
[String::from("default")],
"nobody-serves-this",
None,
[String::from("greet")].iter(),
worker_tx,
)?;
let outcome = parked.join().map_err(|_| "parked dispatch panicked")?;
assert!(outcome.is_err(), "the released dispatch must resolve");
assert!(
state.unserved_queues()?.is_empty(),
"a resolved dispatch must leave the unserved state"
);
registration.deregister()?;
Ok(())
}
#[tokio::test]
async fn namespace_store_is_reachable_and_functional_after_default_boot()
-> Result<(), Box<dyn std::error::Error>> {
use aion_store::{MintOutcome, NamespaceOrigin};
let state =
ServerState::build_with_store(InMemoryStore::default(), runtime_config()).await?;
let store = state.namespace_store();
let outcome = store
.register_namespace("orders", NamespaceOrigin::WorkerMint)
.await?;
assert_eq!(
outcome,
MintOutcome::Created,
"the first reference to a namespace mints it"
);
let again = store
.register_namespace("orders", NamespaceOrigin::WorkerMint)
.await?;
assert_eq!(
again,
MintOutcome::AlreadyExisted,
"a second reference touches the existing record rather than re-creating it"
);
let fetched = store.get_namespace("orders").await?;
let record = fetched.ok_or("registered namespace must be retrievable via get_namespace")?;
assert_eq!(record.name, "orders");
assert_eq!(record.origin, NamespaceOrigin::WorkerMint);
let listed = store.list_namespaces().await?;
assert!(
listed.iter().any(|record| record.name == "orders"),
"list_namespaces returns the minted namespace"
);
Ok(())
}
#[tokio::test(flavor = "multi_thread")]
async fn connect_store_haematite_round_trips_through_event_store()
-> Result<(), Box<dyn std::error::Error>> {
use aion_core::{ContentType, EventEnvelope, PackageVersion, Payload, RunId, WorkflowId};
use aion_store::WriteToken;
use chrono::Utc;
use crate::config::{StoreBackend, StoreConfig};
let data_dir = crate::test_support::private_tempdir()?;
let connected = super::connect_store(StoreConfig {
backend: StoreBackend::Haematite,
owned_shards: Vec::new(),
data_dir: Some(data_dir.path().to_string_lossy().into_owned()),
shard_count: 1,
cluster: None,
node_cache_budget: Some(test_node_cache_budget()?),
..StoreConfig::default()
})
.await?;
let event_store = connected.event_store;
assert!(
connected.outbox_store.is_some(),
"the haematite backend shares its leaf store as the dispatcher's outbox store"
);
assert!(
connected.bootstrap_coordinator,
"a single-node haematite boot owns all shards and bootstraps the coordinator"
);
assert!(
connected.cluster_responder.is_none(),
"a single-node (no [cluster]) haematite boot has no distributed responder"
);
let workflow_id = WorkflowId::new_v4();
let event = aion_core::Event::WorkflowStarted {
envelope: EventEnvelope {
seq: 1,
recorded_at: Utc::now(),
workflow_id: workflow_id.clone(),
},
workflow_type: String::from("checkout"),
input: Payload::new(ContentType::Json, b"{}".to_vec()),
run_id: RunId::new_v4(),
parent_run_id: None,
package_version: PackageVersion::new("a".repeat(64)),
};
event_store
.append(
WriteToken::recorder(),
&workflow_id,
std::slice::from_ref(&event),
0,
)
.await?;
let history = event_store.read_history(&workflow_id).await?;
assert_eq!(
history.len(),
1,
"an event appended through the server's dyn EventStore reads back"
);
Ok(())
}
fn test_node_cache_budget() -> Result<haematite::NodeCacheBudget, Box<dyn std::error::Error>> {
Ok(haematite::NodeCacheBudget::bytes(1 << 30)?)
}
#[cfg(unix)]
#[test]
fn haematite_root_swap_before_first_backend_touch_cannot_redirect_writes()
-> Result<(), Box<dyn std::error::Error>> {
use std::os::unix::fs::symlink;
let sandbox = crate::test_support::private_tempdir()?;
let configured_root = sandbox.path().join("data");
let held_root = sandbox.path().join("held-data");
let outside = sandbox.path().join("outside");
std::fs::create_dir(&outside)?;
let configured = configured_root
.to_str()
.ok_or("temporary data path was not UTF-8")?;
let (store, responder) = super::build_haematite_store_with_hook(
configured,
4,
None,
test_node_cache_budget()?,
|| {
std::fs::rename(&configured_root, &held_root)?;
symlink(&outside, &configured_root)?;
Ok(())
},
)?;
assert!(responder.is_none());
let outside_entries = std::fs::read_dir(&outside)?.collect::<Result<Vec<_>, _>>()?;
assert!(
outside_entries.is_empty(),
"Haematite followed the replaced ambient root and wrote outside"
);
assert!(held_root.join("config.json").is_file());
for shard in 0..4 {
let shard_path = held_root.join(format!("shard-{shard}"));
assert!(shard_path.is_dir(), "shard {shard} was not materialized");
assert!(
std::fs::read_dir(&shard_path)?
.next()
.transpose()?
.is_some(),
"shard {shard} did not run Haematite's materialization path"
);
}
drop(store);
Ok(())
}
#[cfg(any(target_os = "linux", target_os = "android"))]
#[tokio::test]
async fn proc_fd_backend_path_survives_a_post_startup_root_swap()
-> Result<(), Box<dyn std::error::Error>> {
use std::os::unix::fs::symlink;
use aion_core::{ContentType, EventEnvelope, PackageVersion, Payload, RunId, WorkflowId};
use aion_store::{WritableEventStore as _, WriteToken};
use chrono::Utc;
let sandbox = crate::test_support::private_tempdir()?;
let configured_root = sandbox.path().join("data");
let held_root = sandbox.path().join("held-data");
let capture = sandbox.path().join("capture");
std::fs::create_dir(&capture)?;
let configured = configured_root
.to_str()
.ok_or("temporary data path was not UTF-8")?;
let (store, responder) =
super::build_haematite_store(configured, 4, None, test_node_cache_budget()?)?;
assert!(responder.is_none());
std::fs::rename(&configured_root, &held_root)?;
symlink(&capture, &configured_root)?;
let workflow_id = WorkflowId::new_v4();
let event = aion_core::Event::WorkflowStarted {
envelope: EventEnvelope {
seq: 1,
recorded_at: Utc::now(),
workflow_id: workflow_id.clone(),
},
workflow_type: String::from("post-startup-root-swap"),
input: Payload::new(ContentType::Json, b"{}".to_vec()),
run_id: RunId::new_v4(),
parent_run_id: None,
package_version: PackageVersion::new("a".repeat(64)),
};
store
.append(
WriteToken::recorder(),
&workflow_id,
std::slice::from_ref(&event),
0,
)
.await?;
let captured = std::fs::read_dir(&capture)?.collect::<Result<Vec<_>, _>>()?;
assert!(
captured.is_empty(),
"post-startup append followed the replacement symlink into capture"
);
assert!(held_root.join("config.json").is_file());
drop(store);
Ok(())
}
#[cfg(all(unix, not(any(target_os = "linux", target_os = "android"))))]
#[test]
fn path_ambient_haematite_refuses_group_or_world_writable_ancestors()
-> Result<(), Box<dyn std::error::Error>> {
use std::os::unix::fs::PermissionsExt as _;
let sandbox = crate::test_support::private_tempdir()?;
std::fs::set_permissions(sandbox.path(), std::fs::Permissions::from_mode(0o700))?;
for mode in [0o770, 0o1777] {
let shared = sandbox.path().join(format!("shared-{mode:o}"));
let data_root = shared.join("data");
std::fs::create_dir(&shared)?;
std::fs::set_permissions(&shared, std::fs::Permissions::from_mode(mode))?;
std::fs::create_dir(&data_root)?;
std::fs::set_permissions(&data_root, std::fs::Permissions::from_mode(0o700))?;
let configured = data_root
.to_str()
.ok_or("temporary data path was not UTF-8")?;
let Err(error) =
super::build_haematite_store(configured, 4, None, test_node_cache_budget()?)
else {
return Err(format!("mode {mode:04o} ancestor was accepted").into());
};
let message = error.to_string();
let crate::ServerError::UnsafeDataRootAncestor {
data_root: resolved_root,
component,
reason,
} = error
else {
return Err(format!("expected typed unsafe-ancestor error, got {message}").into());
};
assert_eq!(resolved_root, std::fs::canonicalize(&data_root)?);
assert_eq!(component, std::fs::canonicalize(&shared)?);
assert!(
reason.contains(&format!("mode {mode:04o}")),
"unexpected reason: {reason}"
);
if mode & 0o1000 != 0 {
assert!(reason.contains("sticky bit is not accepted"));
}
assert!(message.contains("private Aion home"));
assert!(
!data_root.join("config.json").exists(),
"Haematite touched its ambient path before the refusal"
);
}
Ok(())
}
#[cfg(target_os = "macos")]
#[test]
fn path_ambient_haematite_refuses_mutating_allow_acl_ancestor()
-> Result<(), Box<dyn std::error::Error>> {
use std::os::unix::fs::PermissionsExt as _;
let sandbox = crate::test_support::private_tempdir()?;
std::fs::set_permissions(sandbox.path(), std::fs::Permissions::from_mode(0o700))?;
let shared = sandbox.path().join("acl-shared");
let data_root = shared.join("data");
std::fs::create_dir(&shared)?;
std::fs::set_permissions(&shared, std::fs::Permissions::from_mode(0o700))?;
let acl = "everyone allow list,search,add_file,add_subdirectory,delete_child";
let status = std::process::Command::new("chmod")
.arg("+a")
.arg(acl)
.arg(&shared)
.status()?;
assert!(status.success(), "failed to install Darwin regression ACL");
let configured = data_root
.to_str()
.ok_or("temporary data path was not UTF-8")?;
let result = super::build_haematite_store(configured, 4, None, test_node_cache_budget()?);
let cleanup = std::process::Command::new("chmod")
.arg("-RN")
.arg(&shared)
.status()?;
assert!(cleanup.success(), "failed to clean Darwin regression ACL");
let Err(error) = result else {
return Err("mutating non-euid allow ACL ancestor was accepted".into());
};
let message = error.to_string();
let crate::ServerError::UnsafeDataRootAncestor {
component, reason, ..
} = error
else {
return Err(format!("expected typed unsafe-ancestor error, got {message}").into());
};
assert_eq!(component, std::fs::canonicalize(&shared)?);
assert!(
reason.contains("allow"),
"reason did not name the ACE: {reason}"
);
assert!(
reason.contains("everyone"),
"reason did not name the ACE principal: {reason}"
);
assert!(
!data_root.join("config.json").exists(),
"Haematite touched its ambient path before the ACL refusal"
);
Ok(())
}
#[cfg(target_os = "macos")]
#[test]
fn path_ambient_haematite_accepts_the_euid_uuid_allow_ace()
-> Result<(), Box<dyn std::error::Error>> {
use std::os::unix::fs::PermissionsExt as _;
use exacl::{AclEntry, AclOption, Perm};
let sandbox = crate::test_support::private_tempdir()?;
std::fs::set_permissions(sandbox.path(), std::fs::Permissions::from_mode(0o700))?;
let private_parent = sandbox.path().join("euid-uuid-allow");
let data_root = private_parent.join("data");
std::fs::create_dir(&private_parent)?;
std::fs::set_permissions(&private_parent, std::fs::Permissions::from_mode(0o700))?;
let server_uid = rustix::process::geteuid().as_raw();
let ace_qualifier = crate::filesystem::darwin_user_uuid_for_test(server_uid)?;
let entry = AclEntry::allow_user(
&ace_qualifier.to_string(),
Perm::EXECUTE | Perm::WRITE | Perm::APPEND | Perm::DELETE_CHILD,
None,
);
exacl::setfacl(
&[private_parent.as_path()],
&[entry],
AclOption::SYMLINK_ACL,
)?;
let configured = data_root
.to_str()
.ok_or("temporary data path was not UTF-8")?;
let result = super::build_haematite_store(configured, 4, None, test_node_cache_budget()?);
let cleanup = std::process::Command::new("chmod")
.arg("-RN")
.arg(&private_parent)
.status()?;
assert!(cleanup.success(), "failed to clean euid UUID allow ACL");
let (store, responder) = result?;
assert!(responder.is_none());
assert!(data_root.join("config.json").is_file());
drop(store);
Ok(())
}
#[cfg(target_os = "macos")]
#[test]
fn path_ambient_haematite_refuses_a_non_euid_user_uuid_allow_ace()
-> Result<(), Box<dyn std::error::Error>> {
use std::os::unix::fs::PermissionsExt as _;
use exacl::{AclEntry, AclOption, Perm};
let sandbox = crate::test_support::private_tempdir()?;
std::fs::set_permissions(sandbox.path(), std::fs::Permissions::from_mode(0o700))?;
let shared = sandbox.path().join("non-euid-uuid-allow");
let data_root = shared.join("data");
std::fs::create_dir(&shared)?;
std::fs::set_permissions(&shared, std::fs::Permissions::from_mode(0o700))?;
let server_uid = rustix::process::geteuid().as_raw();
let foreign_uid = u32::from(server_uid == 0);
let foreign_qualifier = crate::filesystem::darwin_user_uuid_for_test(foreign_uid)?;
let entry = AclEntry::allow_user(
&foreign_qualifier.to_string(),
Perm::EXECUTE | Perm::WRITE | Perm::APPEND | Perm::DELETE_CHILD,
None,
);
exacl::setfacl(&[shared.as_path()], &[entry], AclOption::SYMLINK_ACL)?;
let configured = data_root
.to_str()
.ok_or("temporary data path was not UTF-8")?;
let result = super::build_haematite_store(configured, 4, None, test_node_cache_budget()?);
let cleanup = std::process::Command::new("chmod")
.arg("-RN")
.arg(&shared)
.status()?;
assert!(cleanup.success(), "failed to clean non-euid UUID allow ACL");
let Err(error) = result else {
return Err("mutating non-euid user UUID allow ACE was accepted".into());
};
let message = error.to_string();
let crate::ServerError::UnsafeDataRootAncestor {
component, reason, ..
} = error
else {
return Err(format!("expected typed unsafe-ancestor error, got {message}").into());
};
assert_eq!(component, std::fs::canonicalize(&shared)?);
assert!(
reason.contains("allow") && reason.contains(&format!("server euid {server_uid}")),
"reason did not name the rejected ACE: {reason}"
);
assert!(
!data_root.join("config.json").exists(),
"Haematite touched its ambient path before the UUID ACL refusal"
);
Ok(())
}
#[cfg(target_os = "macos")]
#[test]
fn path_ambient_haematite_accepts_a_deny_only_acl_ancestor()
-> Result<(), Box<dyn std::error::Error>> {
use std::os::unix::fs::PermissionsExt as _;
let sandbox = crate::test_support::private_tempdir()?;
std::fs::set_permissions(sandbox.path(), std::fs::Permissions::from_mode(0o700))?;
let private_parent = sandbox.path().join("deny-only");
let data_root = private_parent.join("data");
std::fs::create_dir(&private_parent)?;
std::fs::set_permissions(&private_parent, std::fs::Permissions::from_mode(0o700))?;
let status = std::process::Command::new("chmod")
.arg("+a")
.arg("everyone deny delete")
.arg(&private_parent)
.status()?;
assert!(status.success(), "failed to install Darwin deny-only ACL");
let configured = data_root
.to_str()
.ok_or("temporary data path was not UTF-8")?;
let result = super::build_haematite_store(configured, 4, None, test_node_cache_budget()?);
let cleanup = std::process::Command::new("chmod")
.arg("-RN")
.arg(&private_parent)
.status()?;
assert!(cleanup.success(), "failed to clean Darwin deny-only ACL");
let (store, responder) = result?;
assert!(responder.is_none());
assert!(data_root.join("config.json").is_file());
drop(store);
Ok(())
}
#[cfg(target_os = "macos")]
#[test]
fn path_ambient_haematite_accepts_the_stock_home_acl_chain()
-> Result<(), Box<dyn std::error::Error>> {
use std::os::unix::fs::PermissionsExt as _;
use users::os::unix::UserExt as _;
let effective_uid = rustix::process::geteuid().as_raw();
let effective_user = users::get_user_by_uid(effective_uid)
.ok_or_else(|| format!("server euid {effective_uid} has no account record"))?;
let sandbox = tempfile::Builder::new()
.prefix(".aion-acl-home-proof-")
.tempdir_in(effective_user.home_dir())?;
std::fs::set_permissions(sandbox.path(), std::fs::Permissions::from_mode(0o700))?;
let data_root = sandbox.path().join("data");
let configured = data_root
.to_str()
.ok_or("temporary data path was not UTF-8")?;
let (store, responder) =
super::build_haematite_store(configured, 4, None, test_node_cache_budget()?)?;
assert!(responder.is_none());
assert!(data_root.join("config.json").is_file());
drop(store);
Ok(())
}
#[cfg(all(unix, not(any(target_os = "linux", target_os = "android"))))]
#[test]
fn path_ambient_haematite_accepts_an_owner_controlled_chain()
-> Result<(), Box<dyn std::error::Error>> {
use std::os::unix::fs::PermissionsExt as _;
let sandbox = crate::test_support::private_tempdir()?;
std::fs::set_permissions(sandbox.path(), std::fs::Permissions::from_mode(0o700))?;
let private_parent = sandbox.path().join("private");
let data_root = private_parent.join("data");
std::fs::create_dir(&private_parent)?;
std::fs::set_permissions(&private_parent, std::fs::Permissions::from_mode(0o700))?;
let configured = data_root
.to_str()
.ok_or("temporary data path was not UTF-8")?;
let (store, responder) =
super::build_haematite_store(configured, 4, None, test_node_cache_budget()?)?;
assert!(responder.is_none());
assert!(data_root.join("config.json").is_file());
for shard in 0..4 {
assert!(data_root.join(format!("shard-{shard}")).is_dir());
}
drop(store);
Ok(())
}
#[tokio::test]
async fn connect_store_memory_backend_exposes_no_outbox_store()
-> Result<(), Box<dyn std::error::Error>> {
use crate::config::{StoreBackend, StoreConfig};
let connected = super::connect_store(StoreConfig {
backend: StoreBackend::Memory,
owned_shards: Vec::new(),
data_dir: None,
shard_count: 1,
cluster: None,
node_cache_budget: None,
..StoreConfig::default()
})
.await?;
assert!(
connected.outbox_store.is_none(),
"the in-memory backend exposes no outbox store"
);
Ok(())
}
#[tokio::test]
async fn haematite_boot_refuses_without_a_node_cache_budget()
-> Result<(), Box<dyn std::error::Error>> {
use crate::ServerError;
use crate::config::{StoreBackend, StoreConfig};
let sandbox = crate::test_support::private_tempdir()?;
let data_dir = sandbox.path().join("data");
let error = super::connect_haematite_store(StoreConfig {
backend: StoreBackend::Haematite,
data_dir: Some(
data_dir
.to_str()
.ok_or("temporary data path was not UTF-8")?
.to_owned(),
),
shard_count: 4,
..StoreConfig::default()
})
.await
.err()
.ok_or("the haematite boot path must refuse a store config with no node_cache_budget")?;
let ServerError::Config { message } = error else {
return Err(format!("expected a config refusal, got {error:?}").into());
};
assert!(
message.contains("store.node_cache_budget"),
"the refusal must name the missing key, got: {message}"
);
assert!(
message.contains("AION_STORE_NODE_CACHE_BUDGET"),
"the refusal must name the environment override, got: {message}"
);
Ok(())
}
#[tokio::test]
async fn configured_node_cache_budget_reaches_the_created_database()
-> Result<(), Box<dyn std::error::Error>> {
use crate::config::{StoreBackend, StoreConfig};
const ONE_GIB: usize = 1 << 30;
let sandbox = crate::test_support::private_tempdir()?;
let data_dir = sandbox.path().join("data");
let connected = super::connect_haematite_store(StoreConfig {
backend: StoreBackend::Haematite,
data_dir: Some(
data_dir
.to_str()
.ok_or("temporary data path was not UTF-8")?
.to_owned(),
),
shard_count: 4,
node_cache_budget: Some(haematite::NodeCacheBudget::bytes(ONE_GIB)?),
..StoreConfig::default()
})
.await?;
drop(connected);
let recorded: serde_json::Value =
serde_json::from_slice(&std::fs::read(data_dir.join("config.json"))?)?;
assert_eq!(
recorded.get("node_cache_budget"),
Some(&serde_json::json!({ "bytes": ONE_GIB })),
"the operator's budget must be the one haematite created the database with"
);
Ok(())
}
#[tokio::test]
async fn state_build_fails_without_event_broadcast_capacity()
-> Result<(), Box<dyn std::error::Error>> {
let mut runtime = runtime_config();
runtime.websocket.event_broadcast_capacity = None;
let error = ServerState::build_with_store(InMemoryStore::default(), runtime)
.await
.err()
.ok_or("state build must fail when event streaming is unsized")?;
assert!(error.is_config(), "expected a config error, got {error}");
assert!(
error
.to_string()
.contains("websocket.event_broadcast_capacity"),
"error must name the missing key: {error}"
);
Ok(())
}
#[tokio::test]
async fn state_build_fails_without_query_timeout() -> Result<(), Box<dyn std::error::Error>> {
let mut runtime = runtime_config();
runtime.query_timeout = None;
let error = ServerState::build_with_store(InMemoryStore::default(), runtime)
.await
.err()
.ok_or("state build must fail when the query reply deadline is unset")?;
assert!(error.is_config(), "expected a config error, got {error}");
assert!(
error.to_string().contains("runtime.query_timeout_ms"),
"error must name the missing key: {error}"
);
assert!(
error.to_string().contains("AION_RUNTIME_QUERY_TIMEOUT_MS"),
"error must name the environment override: {error}"
);
Ok(())
}
#[tokio::test]
async fn state_build_fails_with_zero_query_timeout() -> Result<(), Box<dyn std::error::Error>> {
let mut runtime = runtime_config();
runtime.query_timeout = Some(Duration::ZERO);
let error = ServerState::build_with_store(InMemoryStore::default(), runtime)
.await
.err()
.ok_or("state build must fail when the query reply deadline is zero")?;
assert!(error.is_config(), "expected a config error, got {error}");
assert!(
error.to_string().contains("runtime.query_timeout_ms"),
"error must name the zero-valued key: {error}"
);
Ok(())
}
#[tokio::test(flavor = "multi_thread")]
async fn a_completed_check_through_the_built_dispatcher_lands_in_the_returned_slot()
-> Result<(), Box<dyn std::error::Error>> {
use std::collections::{BTreeMap, VecDeque};
use std::sync::{Arc, Mutex};
use aion::ActivityDispatch;
use aion_core::{ActivityId, RunId, WorkflowId};
use aion_package::ActionBodyContract;
use crate::update_check::document::{FETCH_ACTION, FETCH_COMMAND, UPDATE_CHECK_QUEUE};
use crate::worker::{DeclaredBodies, DeclaredBodyLookup, DispatchingRun};
struct SequencedBodies {
replies: Mutex<VecDeque<DeclaredBodyLookup>>,
}
impl DeclaredBodies for SequencedBodies {
fn body_for(
&self,
_task_queue: &str,
_action: &str,
_run: DispatchingRun<'_>,
) -> DeclaredBodyLookup {
let mut replies = match self.replies.lock() {
Ok(replies) => replies,
Err(poisoned) => poisoned.into_inner(),
};
replies.pop_front().unwrap_or(DeclaredBodyLookup::None)
}
}
let runtime = runtime_config();
let cluster_publisher = crate::cluster_publisher::ClusterEventPublisher::new(
ServerState::FALLBACK_CLUSTER_BROADCAST_CAPACITY,
);
let namespace_store: Arc<dyn aion_store::NamespaceStore> =
Arc::new(InMemoryStore::default());
let worker_deployment_store: Arc<dyn aion_store::WorkerDeploymentStore> =
Arc::new(InMemoryStore::default());
let seams = super::build_worker_seams(
&runtime,
&cluster_publisher,
&namespace_store,
&worker_deployment_store,
None,
);
let fixture = concat!(
env!("CARGO_MANIFEST_DIR"),
"/src/update_check/fixtures/aion-cli-index.jsonl"
);
seams.declared_bodies.install(Arc::new(SequencedBodies {
replies: Mutex::new(VecDeque::from([
DeclaredBodyLookup::Declared(ActionBodyContract::Run {
command: FETCH_COMMAND.to_owned(),
}),
DeclaredBodyLookup::Declared(ActionBodyContract::Run {
command: format!("cat {fixture}"),
}),
])),
}));
let transcript = crate::activity_publisher::ActivityEventPublisher::new(
Arc::new(aion_store::InMemoryObservabilityStore::default()),
ServerState::FALLBACK_CLUSTER_BROADCAST_CAPACITY,
crate::activity_publisher::TranscriptBatchPolicy::UNBATCHED,
);
let (dispatcher, _mock_registry, _attempt_owners, _workspace_root, update_status) =
super::build_decorated_dispatcher(&runtime, &seams, transcript);
assert_eq!(
update_status.last(),
None,
"the returned slot must start honestly empty"
);
let dispatch = ActivityDispatch {
namespace: "default".to_owned(),
task_queue: UPDATE_CHECK_QUEUE.to_owned(),
node: None,
workflow_id: WorkflowId::new_v4(),
run_id: RunId::new_v4(),
activity_id: ActivityId::from_sequence_position(1),
name: FETCH_ACTION.to_owned(),
input: "{}".to_owned(),
config: "{}".to_owned(),
attempt: 1,
labels: BTreeMap::new(),
advisory: false,
};
let handle = tokio::task::spawn_blocking(move || dispatcher.dispatch(dispatch));
let encoded = handle
.await?
.map_err(|error| format!("the check dispatch failed: {error}"))?;
let outcome: serde_json::Value = serde_json::from_str(&encoded)?;
assert_eq!(outcome["exit_code"], 0, "the local stand-in command ran");
let recorded = update_status.last().ok_or(
"the completed check must land in the RETURNED slot — the one the boot path \
stores and /update-status serves; an empty slot here is the disconnected-\
producer mis-wire the r1 review proved unmeasured",
)?;
assert_eq!(recorded.latest_known, "0.13.7");
Ok(())
}
}