use std::net::SocketAddr;
use std::process::Command;
use std::sync::Arc;
use std::time::Duration;
use aion::Engine;
use aion_package::{
BeamModule, BeamSet, CURRENT_FORMAT_VERSION, Manifest, ManifestVersion, PackageBuilder,
};
use aion_proto::StreamedEvent;
use aion_server::config::{
AuthConfig, AuthoringConfig, DeployConfig, DevConfig, ListenConfig, MetricsConfig,
NamespaceConfig, NamespaceMode, OpsConsoleAssetSource, OpsConsoleConfig, RuntimeConfig,
WebSocketConfig, WorkerConfig,
};
use aion_server::{ServerState, api::http::http_router};
use aion_store::InMemoryStore;
use axum::body;
use axum::http::{Request, StatusCode, request::Builder};
use futures::{SinkExt, StreamExt};
use serde_json::{Value, json};
use tokio_tungstenite::connect_async;
use tokio_tungstenite::tungstenite::Message;
use tokio_tungstenite::tungstenite::client::IntoClientRequest;
use tower::ServiceExt;
type TestError = Box<dyn std::error::Error>;
const NAMESPACE: &str = "tenant-a";
const FIXTURE_MODULE: &str = "aion_dev_fixture";
const RECEIVE_TIMEOUT: Duration = Duration::from_secs(5);
fn runtime_config(dev_enabled: bool) -> RuntimeConfig {
RuntimeConfig {
listen: ListenConfig {
grpc: SocketAddr::from(([127, 0, 0, 1], 0)),
http: SocketAddr::from(([127, 0, 0, 1], 0)),
},
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),
},
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 {
enabled: dev_enabled,
},
outbox: aion_server::config::OutboxConfig::default(),
observability: aion_server::config::ObservabilityConfig::default(),
scheduler_threads: 1,
query_timeout: Some(Duration::from_secs(10)),
default_namespace: NAMESPACE.to_owned(),
auto_create: aion_server::config::AutoCreate::Open,
max_in_flight_activities: aion_server::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(),
}
}
async fn dev_server(dev_enabled: bool) -> Result<(ServerState, axum::Router), TestError> {
let state =
ServerState::build_with_store(InMemoryStore::default(), runtime_config(dev_enabled))
.await?;
let router = http_router(state.clone())?;
Ok((state, router))
}
fn granted(builder: Builder) -> Builder {
builder
.header("x-aion-subject", "dev")
.header("x-aion-namespaces", NAMESPACE)
}
fn post(path: &str, body: &Value) -> Result<Request<body::Body>, TestError> {
Ok(granted(
Request::builder()
.uri(path)
.method("POST")
.header("content-type", "application/json"),
)
.body(body::Body::from(serde_json::to_vec(body)?))?)
}
async fn read_json(response: axum::response::Response) -> Result<Value, TestError> {
let bytes = body::to_bytes(response.into_body(), usize::MAX).await?;
Ok(serde_json::from_slice(&bytes)?)
}
#[tokio::test]
async fn dev_surface_is_dark_when_disabled() -> Result<(), TestError> {
let (state, router) = dev_server(false).await?;
assert!(
state.activity_mock_registry().is_none(),
"a server with the dev surface disabled installs no mock registry, \
so the engine runs the bare production dispatcher (CN4)"
);
let response = router
.oneshot(post(
"/dev/runs",
&json!({"namespace": NAMESPACE, "workflow_type": "x", "input": {}}),
)?)
.await?;
assert_eq!(
response.status(),
StatusCode::NOT_FOUND,
"every /dev/* path is a plain 404 when the dev surface is dark"
);
Ok(())
}
#[tokio::test]
async fn dev_trigger_drives_the_real_engine_start_path() -> Result<(), TestError> {
let (state, router) = dev_server(true).await?;
assert!(
state.activity_mock_registry().is_some(),
"the dev surface installs the shared mock registry the engine consults"
);
let response = router
.oneshot(post(
"/dev/runs",
&json!({"namespace": NAMESPACE, "workflow_type": "missing", "input": {}}),
)?)
.await?;
assert_eq!(response.status(), StatusCode::NOT_FOUND);
let body = read_json(response).await?;
assert_eq!(body["error_type"], json!("WorkflowTypeNotFound"));
Ok(())
}
#[tokio::test]
async fn dev_replay_rejects_an_unknown_run_over_the_real_store() -> Result<(), TestError> {
let (_state, router) = dev_server(true).await?;
let response = router
.oneshot(post(
"/dev/replay",
&json!({
"namespace": NAMESPACE,
"workflow_id": "00000000-0000-0000-0000-0000000000aa",
}),
)?)
.await?;
assert_eq!(response.status(), StatusCode::NOT_FOUND);
Ok(())
}
#[tokio::test]
async fn dev_mock_registration_requires_the_run_to_exist() -> Result<(), TestError> {
let (_state, router) = dev_server(true).await?;
let response = router
.oneshot(post(
"/dev/mocks",
&json!({
"namespace": NAMESPACE,
"workflow_id": "00000000-0000-0000-0000-0000000000bb",
"activity_name": "charge",
"outcome": {"kind": "succeeds", "result": {"ok": true}},
}),
)?)
.await?;
assert_eq!(response.status(), StatusCode::NOT_FOUND);
Ok(())
}
fn compile_fixture_beam() -> Result<Option<Vec<u8>>, TestError> {
if Command::new("erlc").arg("-version").output().is_err() {
return Ok(None);
}
let dir = std::env::temp_dir().join(format!("aion-dev-ui-{}", uuid::Uuid::new_v4()));
std::fs::create_dir(&dir)?;
let source = dir.join(format!("{FIXTURE_MODULE}.erl"));
let beam = dir.join(format!("{FIXTURE_MODULE}.beam"));
std::fs::write(
&source,
format!(
"-module({FIXTURE_MODULE}).\n\
-export([run/1]).\n\
run(_Input) ->\n\
{{ok, _Release}} = aion_flow_ffi:receive_signal(<<\"release\">>, <<\"{{}}\">>),\n\
1.\n"
),
)?;
let status = Command::new("erlc")
.arg("-o")
.arg(&dir)
.arg(&source)
.status()?;
if !status.success() {
let _ = std::fs::remove_dir_all(&dir);
return Err(format!("erlc failed with status {status}").into());
}
let bytes = std::fs::read(&beam)?;
std::fs::remove_dir_all(&dir)?;
Ok(Some(bytes))
}
fn fixture_archive(beam: &[u8]) -> Result<aion_package::Package, TestError> {
let beams = BeamSet::new(vec![BeamModule::new(FIXTURE_MODULE, beam.to_vec())])?;
let manifest = Manifest {
entry_module: FIXTURE_MODULE.to_owned(),
entry_function: "run".to_owned(),
input_schema: json!({ "type": "object" }),
output_schema: json!({ "type": "integer" }),
timeout: Some(Duration::from_secs(30)),
activities: vec![],
version: ManifestVersion::new("test"),
format_version: CURRENT_FORMAT_VERSION,
additional_workflows: Vec::new(),
};
let bytes = PackageBuilder::new(manifest, beams).write_to_bytes()?;
Ok(aion_package::Package::load_from_bytes(
&bytes,
aion_package::ExtractionLimits::unbounded(),
)?)
}
async fn load_fixture(engine: &Arc<Engine>) -> Result<(), TestError> {
let Some(beam) = compile_fixture_beam()? else {
return Err("erlc absent".into());
};
engine.load_package(fixture_archive(&beam)?).await?;
Ok(())
}
#[tokio::test]
async fn dev_triggered_run_streams_over_the_existing_firehose() -> Result<(), TestError> {
let (state, _router) = dev_server(true).await?;
let engine = state
.deploy_guard()
.engine()
.map(Arc::clone)
.map_err(|error| -> TestError { error.to_string().into() })?;
if load_fixture(&engine).await.is_err() {
tracing::info!("skipping dev firehose e2e: erlc not installed");
println!(
"skipping dev_triggered_run_streams_over_the_existing_firehose: erlc not installed"
);
return Ok(());
}
let router = http_router(state.clone())?;
let listener = tokio::net::TcpListener::bind("127.0.0.1:0").await?;
let address = listener.local_addr()?;
let server = tokio::spawn(async move {
if let Err(error) = axum::serve(listener, router.into_make_service()).await {
tracing::warn!(%error, "dev e2e server exited with error");
}
});
let mut request = format!("ws://{address}/events/stream").into_client_request()?;
request
.headers_mut()
.insert("x-aion-subject", "dev".parse()?);
request
.headers_mut()
.insert("x-aion-namespaces", NAMESPACE.parse()?);
let (mut socket, _response) = connect_async(request).await?;
socket
.send(Message::Text(
json!({
"type": "subscribe",
"subscription": { "firehose": { "namespace": NAMESPACE } }
})
.to_string()
.into(),
))
.await?;
tokio::time::sleep(Duration::from_millis(200)).await;
let trigger = dev_post(
&state,
"/dev/runs",
&json!({"namespace": NAMESPACE, "workflow_type": FIXTURE_MODULE, "input": {}}),
)
.await?;
let workflow_id = trigger["workflow_id"]
.as_str()
.ok_or("trigger response missing workflow id")?
.to_owned();
assert_eq!(
trigger["stream_subscription"]["path"], "/events/stream",
"the dev trigger must hand back the EXISTING firehose path, not a new stream"
);
let started = tokio::time::timeout(RECEIVE_TIMEOUT, next_event_for(&mut socket, &workflow_id))
.await
.map_err(|_| "timed out waiting for the streamed WorkflowStarted event")??;
assert_eq!(
started.decode_event()?.workflow_id().to_string(),
workflow_id,
"the firehose must deliver the triggered run's own events"
);
server.abort();
Ok(())
}
async fn dev_post(state: &ServerState, path: &str, body: &Value) -> Result<Value, TestError> {
let router = http_router(state.clone())?;
let response = router.oneshot(post(path, body)?).await?;
let status = response.status();
let value = read_json(response).await?;
if !status.is_success() {
return Err(format!("{path} failed with {status}: {value}").into());
}
Ok(value)
}
async fn next_event_for(
socket: &mut tokio_tungstenite::WebSocketStream<
tokio_tungstenite::MaybeTlsStream<tokio::net::TcpStream>,
>,
workflow_id: &str,
) -> Result<StreamedEvent, TestError> {
loop {
let Some(frame) = socket.next().await else {
return Err("socket closed before delivering the run's event".into());
};
let Message::Text(text) = frame? else {
continue;
};
let streamed: StreamedEvent = serde_json::from_str(&text)?;
if streamed.decode_event()?.workflow_id().to_string() == workflow_id {
return Ok(streamed);
}
}
}