use std::time::Duration;
use acton_macro::{acton_actor, acton_message};
use acton_reactive::ipc::{socket_exists, IpcConfig};
use acton_reactive::prelude::*;
use tracing_subscriber::EnvFilter;
#[acton_message(ipc)]
struct Add {
a: i64,
b: i64,
}
#[acton_message(ipc)]
struct Multiply {
a: i64,
b: i64,
}
#[acton_message(ipc)]
struct CalcResult {
result: i64,
operation: String,
}
#[acton_message(ipc)]
struct SearchQuery {
query: String,
limit: usize,
}
#[acton_message(ipc)]
struct SearchResult {
id: usize,
title: String,
score: f64,
}
#[acton_message(ipc)]
struct LogEvent {
level: String,
message: String,
}
#[acton_message(ipc)]
struct PriceUpdate {
symbol: String,
price: f64,
change: f64,
}
#[acton_message(ipc)]
struct StatusChange {
service: String,
status: String,
timestamp_ms: u64,
}
#[acton_actor]
struct CalculatorState {
operations_performed: usize,
}
#[acton_actor]
struct SearchState {
searches_performed: usize,
}
#[acton_actor]
struct LoggerState {
log_count: usize,
}
#[acton_actor]
struct PricePublisherState {
tick_count: usize,
}
#[acton_message]
struct PublishTick;
async fn create_calculator_actor(runtime: &mut ActorRuntime) -> ActorHandle {
let mut calculator = runtime.new_actor_with_name::<CalculatorState>("calculator".to_string());
calculator.mutate_on::<Add>(|actor, envelope| {
let msg = envelope.message();
let result = msg.a + msg.b;
actor.model.operations_performed += 1;
let response = CalcResult {
result,
operation: format!("{} + {}", msg.a, msg.b),
};
println!(
" [Calculator] Add: {} + {} = {} (op #{})",
msg.a, msg.b, result, actor.model.operations_performed
);
let reply_envelope = envelope.reply_envelope();
Reply::pending(async move {
reply_envelope.send(response).await;
})
});
calculator.mutate_on::<Multiply>(|actor, envelope| {
let msg = envelope.message();
let result = msg.a * msg.b;
actor.model.operations_performed += 1;
let response = CalcResult {
result,
operation: format!("{} × {}", msg.a, msg.b),
};
println!(
" [Calculator] Multiply: {} × {} = {} (op #{})",
msg.a, msg.b, result, actor.model.operations_performed
);
let reply_envelope = envelope.reply_envelope();
Reply::pending(async move {
reply_envelope.send(response).await;
})
});
calculator.start().await
}
async fn create_search_actor(runtime: &mut ActorRuntime) -> ActorHandle {
let mut search = runtime.new_actor_with_name::<SearchState>("search".to_string());
search.mutate_on::<SearchQuery>(|actor, envelope| {
let msg = envelope.message();
let query = msg.query.clone();
let limit = msg.limit.clamp(1, 10); actor.model.searches_performed += 1;
println!(
" [Search] Query: \"{}\" (limit: {}, search #{})",
query, limit, actor.model.searches_performed
);
let reply_envelope = envelope.reply_envelope();
Reply::pending(async move {
let sample_items = [
"Getting Started Guide",
"API Reference",
"Tutorial: Building Actors",
"Configuration Options",
"Performance Tuning",
"Troubleshooting FAQ",
"Architecture Overview",
"Migration Guide",
"Security Best Practices",
"Release Notes",
];
for (idx, title) in sample_items.iter().take(limit).enumerate() {
#[allow(clippy::cast_precision_loss)]
let score = (idx as f64).mul_add(-0.1, 1.0);
let result = SearchResult {
id: idx + 1,
title: format!("{title} (matches: {query})"),
score,
};
println!(" [Search] Sending result {}/{}", idx + 1, limit);
reply_envelope.send(result).await;
tokio::time::sleep(Duration::from_millis(100)).await;
}
println!(" [Search] Stream complete");
})
});
search.start().await
}
async fn create_logger_actor(runtime: &mut ActorRuntime) -> ActorHandle {
let mut logger = runtime.new_actor_with_name::<LoggerState>("logger".to_string());
logger.mutate_on::<LogEvent>(|actor, envelope| {
let msg = envelope.message();
actor.model.log_count += 1;
println!(
" [Logger] #{} [{}] {}",
actor.model.log_count,
msg.level.to_uppercase(),
msg.message
);
Reply::ready()
});
logger.start().await
}
async fn create_price_publisher(runtime: &mut ActorRuntime) -> ActorHandle {
let mut publisher =
runtime.new_actor_with_name::<PricePublisherState>("price_publisher".to_string());
publisher.mutate_on::<PublishTick>(|actor, _envelope| {
actor.model.tick_count += 1;
let tick = actor.model.tick_count;
let symbols = ["AAPL", "GOOGL", "MSFT", "AMZN"];
let symbol = symbols[tick % symbols.len()];
let base_price = match symbol {
"AAPL" => 150.0,
"GOOGL" => 140.0,
"MSFT" => 370.0,
"AMZN" => 180.0,
_ => 100.0,
};
#[allow(clippy::cast_precision_loss)]
let change = f64::mul_add((tick as f64 * 0.7).sin(), 2.0, 0.0);
let change = (change * 100.0).round() / 100.0;
let price = base_price + change;
let update = PriceUpdate {
symbol: symbol.to_string(),
price,
change,
};
println!(" [Publisher] Publishing: {symbol} @ ${price:.2} ({change:+.2})");
let broker = actor.broker().clone();
Reply::pending(async move {
broker.broadcast(update).await;
})
});
publisher.start().await
}
fn spawn_price_ticker(publisher: ActorHandle) {
tokio::spawn(async move {
let mut interval = tokio::time::interval(Duration::from_secs(5));
loop {
interval.tick().await;
publisher.send(PublishTick).await;
}
});
}
#[acton_main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
tracing_subscriber::fmt()
.with_env_filter(EnvFilter::from_default_env().add_directive("acton=info".parse()?))
.init();
println!("╔══════════════════════════════════════════════════════════════╗");
println!("║ IPC Client Libraries Example Server ║");
println!("╚══════════════════════════════════════════════════════════════╝");
println!();
let mut runtime = ActonApp::launch_async().await;
let registry = runtime.ipc_registry();
registry.register::<Add>("Add");
registry.register::<Multiply>("Multiply");
registry.register::<CalcResult>("CalcResult");
registry.register::<SearchQuery>("SearchQuery");
registry.register::<SearchResult>("SearchResult");
registry.register::<LogEvent>("LogEvent");
registry.register::<PriceUpdate>("PriceUpdate");
registry.register::<StatusChange>("StatusChange");
println!("📝 Registered {} IPC message types", registry.len());
let calculator = create_calculator_actor(&mut runtime).await;
println!("🧮 Calculator service started");
let search = create_search_actor(&mut runtime).await;
println!("🔍 Search service started");
let logger = create_logger_actor(&mut runtime).await;
println!("📋 Logger service started");
let price_publisher = create_price_publisher(&mut runtime).await;
println!("💰 Price publisher started");
runtime.ipc_expose("calculator", calculator.clone());
runtime.ipc_expose("search", search.clone());
runtime.ipc_expose("logger", logger.clone());
runtime.ipc_expose("price_publisher", price_publisher.clone());
println!("🔗 Exposed actors: calculator, search, logger, price_publisher");
spawn_price_ticker(price_publisher);
println!("⏰ Price ticker started (every 5 seconds)");
let mut ipc_config = IpcConfig::load();
ipc_config.socket.app_name = Some("ipc_client_example".to_string());
let socket_path = ipc_config.socket_path();
if let Some(parent) = socket_path.parent() {
std::fs::create_dir_all(parent)?;
}
let listener_handle = runtime.start_ipc_listener_with_config(ipc_config).await?;
println!("🚀 IPC listener started");
tokio::time::sleep(Duration::from_millis(50)).await;
if socket_exists(&socket_path) {
println!("📡 Socket ready: {}", socket_path.display());
}
println!();
println!("════════════════════════════════════════════════════════════════");
println!(" Server is ready! Test with the client libraries:");
println!();
println!(" Python:");
println!(" cd examples/ipc_client_libraries/python");
println!(" python example_client.py");
println!();
println!(" Node.js:");
println!(" cd examples/ipc_client_libraries/nodejs");
println!(" npm install && npx ts-node src/example-client.ts");
println!();
println!(" Available services:");
println!(" - calculator: Add {{ a, b }}, Multiply {{ a, b }}");
println!(" - search: SearchQuery {{ query, limit }} (streaming)");
println!(" - logger: LogEvent {{ level, message }} (fire-and-forget)");
println!(" - Subscribe to: PriceUpdate, StatusChange");
println!("════════════════════════════════════════════════════════════════");
println!();
println!("Press Ctrl+C to shutdown...");
println!();
tokio::signal::ctrl_c().await?;
println!();
println!("Shutting down...");
listener_handle.stop();
tokio::time::sleep(Duration::from_millis(100)).await;
runtime.shutdown_all().await?;
println!("Server shutdown complete.");
Ok(())
}