use std::path::PathBuf;
use std::time::Duration;
use acton_reactive::ipc::protocol::{is_stream, read_frame, write_envelope, MAX_FRAME_SIZE};
use acton_reactive::ipc::{socket_exists, socket_is_alive, IpcConfig, IpcEnvelope, IpcStreamFrame};
use acton_reactive::prelude::acton_main;
use tokio::net::UnixStream;
use tokio::time::timeout;
const DEFAULT_SERVER: &str = "ipc_streaming_server";
fn resolve_socket_path() -> PathBuf {
let args: Vec<String> = std::env::args().collect();
if let Some(pos) = args.iter().position(|a| a == "--server") {
if let Some(server_name) = args.get(pos + 1) {
let mut config = IpcConfig::load();
config.socket.app_name = Some(server_name.clone());
return config.socket_path();
}
}
if let Some(arg) = args.get(1) {
if !arg.starts_with("--") {
return PathBuf::from(arg);
}
}
let mut config = IpcConfig::load();
config.socket.app_name = Some(DEFAULT_SERVER.to_string());
config.socket_path()
}
async fn send_stream_request(
reader: &mut tokio::net::unix::OwnedReadHalf,
writer: &mut tokio::net::unix::OwnedWriteHalf,
target: &str,
message_type: &str,
payload: serde_json::Value,
timeout_ms: u64,
) -> Result<Vec<IpcStreamFrame>, Box<dyn std::error::Error>> {
let envelope =
IpcEnvelope::new_stream_request_with_timeout(target, message_type, payload, timeout_ms);
println!("📤 Stream request to {target}::{message_type}");
println!(" correlation_id: {}", envelope.correlation_id);
println!(" expects_stream: {}", envelope.expects_stream);
println!(" timeout: {}ms", envelope.response_timeout_ms);
println!();
write_envelope(writer, &envelope).await?;
let mut frames = Vec::new();
let overall_timeout = Duration::from_millis(timeout_ms + 5000);
let frame_timeout = Duration::from_millis(timeout_ms.max(1000));
let start = std::time::Instant::now();
loop {
if start.elapsed() > overall_timeout {
return Err("Stream timeout: no final frame received".into());
}
let result = timeout(frame_timeout, read_frame(reader, MAX_FRAME_SIZE)).await;
match result {
Ok(Ok((msg_type, _format, payload))) => {
if !is_stream(msg_type) {
println!("⚠️ Unexpected message type: 0x{msg_type:02X}");
continue;
}
let frame: IpcStreamFrame = serde_json::from_slice(&payload)?;
display_stream_frame(&frame);
let is_final = frame.is_final;
frames.push(frame);
if is_final {
println!("\n ✅ Stream complete ({} frames received)", frames.len());
break;
}
}
Ok(Err(e)) => {
return Err(format!("Error reading frame: {e}").into());
}
Err(_) => {
return Err("Frame timeout: no response received".into());
}
}
}
Ok(frames)
}
fn display_stream_frame(frame: &IpcStreamFrame) {
if let Some(error) = &frame.error {
println!(" 📥 Frame #{}: ERROR - {}", frame.sequence, error);
} else if let Some(payload) = &frame.payload {
let final_marker = if frame.is_final { " [FINAL]" } else { "" };
if let Some(obj) = payload.as_object() {
println!(" 📥 Frame #{}:{}", frame.sequence, final_marker);
for (key, value) in obj {
println!(" {key}: {value}");
}
} else {
println!(
" 📥 Frame #{}: {}{}",
frame.sequence,
serde_json::to_string(payload).unwrap_or_default(),
final_marker
);
}
} else {
println!(" 📥 Frame #{}: (empty payload)", frame.sequence);
}
}
async fn demo_countdown(
reader: &mut tokio::net::unix::OwnedReadHalf,
writer: &mut tokio::net::unix::OwnedWriteHalf,
) -> Result<(), Box<dyn std::error::Error>> {
println!("╔══════════════════════════════════════════════════════════════╗");
println!("║ Countdown Stream Demo ║");
println!("╚══════════════════════════════════════════════════════════════╝");
println!();
println!("⏱️ Test 1: Countdown from 5 (500ms delay)");
let _frames = send_stream_request(
reader,
writer,
"countdown",
"CountdownRequest",
serde_json::json!({ "start": 5, "delay_ms": 500 }),
10000,
)
.await?;
println!();
println!("⏱️ Test 2: Quick countdown from 3 (100ms delay)");
let _frames = send_stream_request(
reader,
writer,
"countdown",
"CountdownRequest",
serde_json::json!({ "start": 3, "delay_ms": 100 }),
5000,
)
.await?;
Ok(())
}
async fn demo_paginated_list(
reader: &mut tokio::net::unix::OwnedReadHalf,
writer: &mut tokio::net::unix::OwnedWriteHalf,
) -> Result<(), Box<dyn std::error::Error>> {
println!("\n╔══════════════════════════════════════════════════════════════╗");
println!("║ Paginated List Stream Demo ║");
println!("╚══════════════════════════════════════════════════════════════╝");
println!();
println!("📋 Test 1: List items (page size: 3)");
let _frames = send_stream_request(
reader,
writer,
"list_service",
"ListItemsRequest",
serde_json::json!({ "page_size": 3 }),
10000,
)
.await?;
println!();
println!("📋 Test 2: List items (page size: 5)");
let _frames = send_stream_request(
reader,
writer,
"list_service",
"ListItemsRequest",
serde_json::json!({ "page_size": 5 }),
10000,
)
.await?;
Ok(())
}
async fn demo_error_handling(
reader: &mut tokio::net::unix::OwnedReadHalf,
writer: &mut tokio::net::unix::OwnedWriteHalf,
) -> Result<(), Box<dyn std::error::Error>> {
println!("\n╔══════════════════════════════════════════════════════════════╗");
println!("║ Error Handling Demo ║");
println!("╚══════════════════════════════════════════════════════════════╝");
println!();
println!("❌ Test: Stream request to non-existent actor 'nonexistent'");
let result = send_stream_request(
reader,
writer,
"nonexistent",
"SomeMessage",
serde_json::json!({ "test": "data" }),
5000,
)
.await;
match result {
Ok(frames) => {
if let Some(frame) = frames.first() {
if frame.error.is_some() {
println!(" Expected error received");
}
}
}
Err(e) => {
println!(" Error: {e}");
}
}
Ok(())
}
#[acton_main]
async fn main() -> Result<(), Box<dyn std::error::Error>> {
println!("╔══════════════════════════════════════════════════════════════╗");
println!("║ IPC Streaming Response Example (Client) ║");
println!("╚══════════════════════════════════════════════════════════════╝");
println!();
let socket_path = resolve_socket_path();
println!("🔌 Socket path: {}", socket_path.display());
if !socket_exists(&socket_path) {
eprintln!(
"\n❌ Error: Socket does not exist at {}",
socket_path.display()
);
eprintln!(" Make sure the ipc_streaming server is running:");
eprintln!(" cargo run --example ipc_streaming_server --features ipc");
std::process::exit(1);
}
if !socket_is_alive(&socket_path).await {
eprintln!("\n❌ Error: Socket exists but is not responding");
eprintln!(" The server may have crashed. Try restarting it.");
std::process::exit(1);
}
println!("🔗 Connecting to server...");
let stream = UnixStream::connect(&socket_path).await?;
println!("✅ Connected successfully!");
println!();
let (mut reader, mut writer) = stream.into_split();
demo_countdown(&mut reader, &mut writer).await?;
demo_paginated_list(&mut reader, &mut writer).await?;
demo_error_handling(&mut reader, &mut writer).await?;
println!("\n════════════════════════════════════════════════════════════════");
println!(" All streaming IPC tests completed!");
println!("════════════════════════════════════════════════════════════════");
Ok(())
}