#[path = "ipc_server/slot.rs"]
mod slot;
use self::slot::{SlotGuard, SlotPool, MAX_IPC_SLOTS};
use super::arbitrator::{Arbitrator, Outcome};
use super::calibration::Calibration;
use super::config::IpcConfig;
use super::gimbal_handle::GimbalHandle;
use super::models::{CommandMode, ControlSource, GimbalCommand, IpcCommand, IpcResponse, PanMode};
use super::state::StateManager;
use std::sync::Arc;
use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader};
use tokio::net::{UnixListener, UnixStream};
use tracing::{debug, error, info, warn};
pub struct IpcServer {
config: IpcConfig,
state_manager: StateManager,
arbitrator: Arc<Arbitrator>,
gimbal: GimbalHandle,
slots: SlotPool,
}
impl IpcServer {
pub fn new(config: IpcConfig, state_manager: StateManager, gimbal: GimbalHandle) -> Self {
let arbitrator = Arc::new(Arbitrator::new(state_manager.clone()));
let slots = SlotPool::new(state_manager.clone());
Self {
config,
state_manager,
arbitrator,
gimbal,
slots,
}
}
pub async fn start(self) -> crate::error::Result<()> {
let socket_path = &self.config.socket_path;
if std::path::Path::new(socket_path).exists() {
std::fs::remove_file(socket_path)?;
}
info!("Starting IPC server on {}", socket_path);
let listener = UnixListener::bind(socket_path)?;
loop {
match listener.accept().await {
Ok((stream, _addr)) => {
let Some(slot_guard) = self.slots.acquire() else {
warn!(
"IPC server at capacity ({} clients); rejecting new connection",
MAX_IPC_SLOTS
);
drop(stream);
continue;
};
let client_id = slot_guard.id();
info!("New IPC client connected: ID {}", client_id);
let state_manager = self.state_manager.clone();
let arbitrator = self.arbitrator.clone();
let gimbal = self.gimbal.clone();
tokio::spawn(async move {
if let Err(e) = Self::handle_client(
slot_guard,
stream,
state_manager,
arbitrator,
gimbal,
)
.await
{
error!("Client {} error: {}", client_id, e);
}
info!("Client {} disconnected", client_id);
});
}
Err(e) => {
error!("Failed to accept connection: {}", e);
}
}
}
}
async fn handle_client(
slot_guard: SlotGuard,
stream: UnixStream,
state_manager: StateManager,
arbitrator: Arc<Arbitrator>,
gimbal: GimbalHandle,
) -> crate::error::Result<()> {
let client_id = slot_guard.id();
let (reader, mut writer) = stream.into_split();
let mut reader = BufReader::new(reader);
let mut line = String::new();
loop {
line.clear();
match reader.read_line(&mut line).await {
Ok(0) => break, Ok(_) => {
let trimmed = line.trim();
if trimmed.is_empty() {
continue;
}
debug!("Client {}: received command: {}", client_id, trimmed);
let response = Self::process_command(
client_id,
trimmed,
&state_manager,
&arbitrator,
&gimbal,
)
.await;
let response_json = serde_json::to_string(&response)? + "\n";
writer.write_all(response_json.as_bytes()).await?;
writer.flush().await?;
}
Err(e) => {
error!("Client {} read error: {}", client_id, e);
break;
}
}
}
Ok(())
}
async fn process_command(
client_id: u32,
command_str: &str,
state_manager: &StateManager,
arbitrator: &Arbitrator,
gimbal: &GimbalHandle,
) -> IpcResponse {
let command: IpcCommand = match serde_json::from_str(command_str) {
Ok(cmd) => cmd,
Err(e) => return IpcResponse::error(format!("Invalid JSON: {}", e)),
};
debug!("Client {}: parsed command: {:?}", client_id, command);
match command {
IpcCommand::Set { pitch, roll, yaw } => {
let gimbal_cmd = GimbalCommand::new(
ControlSource::UnixSocket(client_id),
CommandMode::Position,
Some(yaw),
Some(pitch),
Some(roll),
);
execute_position(arbitrator, gimbal, gimbal_cmd, "set").await
}
IpcCommand::Center => {
let gimbal_cmd = GimbalCommand::new(
ControlSource::UnixSocket(client_id),
CommandMode::Position,
Some(0.0),
Some(0.0),
Some(0.0),
);
execute_position(arbitrator, gimbal, gimbal_cmd, "center").await
}
IpcCommand::Status => {
let state = state_manager.get_state();
let data = serde_json::json!({
"yaw": state.yaw,
"pitch": state.pitch,
"roll": state.roll,
"pan_mode": state.pan_mode.to_name(),
"standby": state.standby,
"firmware_version": state.firmware_version,
"yaw_offset_deg": gimbal.yaw_offset_deg(),
});
IpcResponse::ok_with_data(data)
}
IpcCommand::Version => match gimbal.get_version().await {
Ok(version) => {
let data = serde_json::json!({ "version": version });
IpcResponse::ok_with_data(data)
}
Err(e) => IpcResponse::error(format!("Failed to get version: {}", e)),
},
IpcCommand::PanMode { mode } => {
if let Some(pan_mode) = PanMode::from_u8(mode) {
match gimbal.set_pan_mode(mode).await {
Ok(_) => {
state_manager.update_pan_mode(pan_mode);
IpcResponse::ok()
}
Err(e) => IpcResponse::error(format!("Failed to set pan mode: {}", e)),
}
} else {
IpcResponse::error(format!("Invalid pan mode: {}", mode))
}
}
IpcCommand::Standby { enabled } => match gimbal.set_standby(enabled).await {
Ok(_) => {
state_manager.update_standby(enabled);
IpcResponse::ok()
}
Err(e) => IpcResponse::error(format!("Failed to set standby: {}", e)),
},
IpcCommand::CalibrateYaw => calibrate_yaw(gimbal).await,
IpcCommand::SetYawOffset { deg } => set_yaw_offset(gimbal, deg),
IpcCommand::Selftest => selftest(gimbal, state_manager).await,
IpcCommand::Help => {
let help_text = serde_json::json!({
"commands": [
{"cmd": "set", "params": {"pitch": "f32", "roll": "f32", "yaw": "f32"}},
{"cmd": "center"},
{"cmd": "status"},
{"cmd": "version"},
{"cmd": "pan_mode", "params": {"mode": "u8"}},
{"cmd": "standby", "params": {"enabled": "bool"}},
{"cmd": "calibrate_yaw"},
{"cmd": "set_yaw_offset", "params": {"deg": "f32"}},
{"cmd": "selftest"},
{"cmd": "help"},
]
});
IpcResponse::ok_with_data(help_text)
}
}
}
}
const CALIBRATION_SETTLE: std::time::Duration = std::time::Duration::from_millis(2500);
async fn calibrate_yaw(gimbal: &GimbalHandle) -> IpcResponse {
let prior_offset = gimbal.yaw_offset_deg();
gimbal.set_yaw_offset_deg(0.0);
if let Err(e) = gimbal.set_attitude(0.0, 0.0, 0.0).await {
gimbal.set_yaw_offset_deg(prior_offset);
return IpcResponse::error(format!("Calibration failed sending center command: {}", e));
}
tokio::time::sleep(CALIBRATION_SETTLE).await;
let raw = match gimbal.get_attitude_raw().await {
Ok(att) => att,
Err(e) => {
gimbal.set_yaw_offset_deg(prior_offset);
return IpcResponse::error(format!("Calibration failed reading attitude: {}", e));
}
};
if !raw.yaw.is_finite() {
gimbal.set_yaw_offset_deg(prior_offset);
return IpcResponse::error(format!(
"Calibration read non-finite yaw: {} — leaving prior offset {:.3}° intact",
raw.yaw, prior_offset
));
}
let new_offset = raw.yaw;
gimbal.set_yaw_offset_deg(new_offset);
let mut persistence_warning = None;
if let Err(e) = (Calibration {
yaw_offset_deg: new_offset,
})
.save()
{
warn!(
"Calibration succeeded in-memory but failed to persist: {} \
— offset will be lost on daemon restart",
e
);
persistence_warning = Some(format!(
"calibration applied in memory but persistence failed: {} \
(will not survive daemon restart)",
e
));
}
info!(
"Yaw calibration: prior_offset={:.3}° → new_offset={:.3}° (raw IMU yaw at SP=0)",
prior_offset, new_offset
);
let mut data = serde_json::json!({
"yaw_offset_deg": new_offset,
"prior_yaw_offset_deg": prior_offset,
});
if let Some(w) = persistence_warning {
data["warning"] = serde_json::Value::String(w);
}
IpcResponse::ok_with_data(data)
}
fn set_yaw_offset(gimbal: &GimbalHandle, deg: f32) -> IpcResponse {
if !deg.is_finite() {
return IpcResponse::error(format!("yaw offset must be finite (got {deg})"));
}
gimbal.set_yaw_offset_deg(deg);
let mut persistence_warning = None;
if let Err(e) = (Calibration {
yaw_offset_deg: deg,
})
.save()
{
warn!(
"set_yaw_offset succeeded in-memory but failed to persist: {} \
— offset will be lost on daemon restart",
e
);
persistence_warning = Some(format!(
"offset applied in memory but persistence failed: {} \
(will not survive daemon restart)",
e
));
}
info!("Yaw offset set explicitly: {:.3}°", deg);
let mut data = serde_json::json!({ "yaw_offset_deg": deg });
if let Some(w) = persistence_warning {
data["warning"] = serde_json::Value::String(w);
}
IpcResponse::ok_with_data(data)
}
const SELFTEST_TOLERANCE_DEG: f32 = 1.0;
async fn selftest(gimbal: &GimbalHandle, state_manager: &StateManager) -> IpcResponse {
if state_manager.is_standby() {
return IpcResponse::error(
"selftest: gimbal is in standby (motors disengaged). \
Toggle standby off before running selftest.",
);
}
let original = match gimbal.get_attitude().await {
Ok(att) => att,
Err(e) => {
return IpcResponse::error(format!("selftest: failed to read initial pose: {e}"));
}
};
let test_points: &[(f32, f32, f32)] = &[
(0.0, 0.0, 0.0),
(15.0, 0.0, 0.0),
(-15.0, 0.0, 0.0),
(0.0, 0.0, 15.0),
(0.0, 0.0, -15.0),
(0.0, 0.0, 0.0),
];
let mut samples: Vec<serde_json::Value> = Vec::with_capacity(test_points.len());
let mut max_pitch_err = 0.0_f32;
let mut max_roll_err = 0.0_f32;
let mut max_yaw_err = 0.0_f32;
let mut passed = true;
let mut aborted_at: Option<String> = None;
for &(p, r, y) in test_points {
if let Err(e) = gimbal.set_attitude(p, r, y).await {
aborted_at = Some(format!("set_attitude({p}, {r}, {y}) failed: {e}"));
passed = false;
break;
}
tokio::time::sleep(CALIBRATION_SETTLE).await;
let pv = match gimbal.get_attitude().await {
Ok(att) => att,
Err(e) => {
aborted_at = Some(format!("get_attitude after SP=({p}, {r}, {y}) failed: {e}"));
passed = false;
break;
}
};
let ep = (p - pv.pitch).abs();
let er = (r - pv.roll).abs();
let ey = (y - pv.yaw).abs();
max_pitch_err = max_pitch_err.max(ep);
max_roll_err = max_roll_err.max(er);
max_yaw_err = max_yaw_err.max(ey);
if ep > SELFTEST_TOLERANCE_DEG || er > SELFTEST_TOLERANCE_DEG || ey > SELFTEST_TOLERANCE_DEG
{
passed = false;
}
samples.push(serde_json::json!({
"sp": {"pitch": p, "roll": r, "yaw": y},
"pv": {"pitch": pv.pitch, "roll": pv.roll, "yaw": pv.yaw},
"error_deg": {"pitch": ep, "roll": er, "yaw": ey},
}));
}
let restore_warning = match gimbal
.set_attitude(original.pitch, original.roll, original.yaw)
.await
{
Ok(()) => None,
Err(e) => {
warn!("selftest: failed to restore initial pose: {e}");
Some(format!(
"selftest sweep finished but the initial pose could not be restored: {e}"
))
}
};
info!(
"Selftest: passed={}, max_err deg pitch={:.3} roll={:.3} yaw={:.3}",
passed, max_pitch_err, max_roll_err, max_yaw_err
);
let mut data = serde_json::json!({
"passed": passed,
"tolerance_deg": SELFTEST_TOLERANCE_DEG,
"samples": samples,
"max_error_deg": {
"pitch": max_pitch_err,
"roll": max_roll_err,
"yaw": max_yaw_err,
},
});
if let Some(w) = aborted_at {
data["aborted_at"] = serde_json::Value::String(w);
}
if let Some(w) = restore_warning {
data["warning"] = serde_json::Value::String(w);
}
IpcResponse::ok_with_data(data)
}
async fn execute_position(
arbitrator: &Arbitrator,
gimbal: &GimbalHandle,
cmd: GimbalCommand,
label: &str,
) -> IpcResponse {
match arbitrator.arbitrate(cmd) {
Outcome::Execute { pitch, roll, yaw } => {
match gimbal.set_attitude(pitch, roll, yaw).await {
Ok(()) => IpcResponse::ok(),
Err(e) => IpcResponse::error(format!("Failed to {label}: {e}")),
}
}
Outcome::Rejected => IpcResponse::error("Command rejected (lower priority)"),
}
}