#[path = "mavlink_manager/handlers.rs"]
mod handlers;
#[path = "mavlink_manager/math.rs"]
mod math;
use super::arbitrator::Arbitrator;
use super::config::{MavlinkConfig, Transport};
use super::gimbal_handle::GimbalHandle;
use super::state::StateManager;
use std::collections::HashSet;
use std::net::SocketAddr;
use std::sync::atomic::AtomicU8;
use std::sync::Arc;
use std::time::Instant;
use tokio::net::UdpSocket;
use tokio::sync::{Mutex, Notify};
use tracing::{debug, error, info};
pub(crate) type Peers = Arc<Mutex<HashSet<SocketAddr>>>;
pub(crate) struct RxCtx<'a> {
pub socket: &'a Arc<UdpSocket>,
pub config: &'a MavlinkConfig,
pub seq: &'a Arc<AtomicU8>,
pub state: &'a StateManager,
pub arbitrator: &'a Arbitrator,
pub gimbal: &'a GimbalHandle,
pub start_time: Instant,
}
pub struct MavlinkManager {
state_manager: StateManager,
arbitrator: Arc<Arbitrator>,
gimbal: GimbalHandle,
peers: Peers,
seq: Arc<AtomicU8>,
start_time: Instant,
notify_lost: Arc<Notify>,
}
impl MavlinkManager {
pub fn new(
state_manager: StateManager,
gimbal: GimbalHandle,
notify_lost: Arc<Notify>,
) -> Self {
let arbitrator = Arc::new(Arbitrator::new(state_manager.clone()));
Self {
state_manager,
arbitrator,
gimbal,
peers: Arc::new(Mutex::new(HashSet::new())),
seq: Arc::new(AtomicU8::new(0)),
start_time: Instant::now(),
notify_lost,
}
}
pub async fn start(self, config: MavlinkConfig) -> crate::error::Result<()> {
match config.transport {
Transport::Udp => self.start_udp(config).await,
}
}
async fn start_udp(self, config: MavlinkConfig) -> crate::error::Result<()> {
let bind_addr = SocketAddr::new(config.bind_addr, config.udp_port);
info!("Starting MAVLink manager on UDP {}", bind_addr);
let socket = Arc::new(UdpSocket::bind(bind_addr).await?);
info!("MAVLink UDP socket bound to {}", bind_addr);
{
let socket = socket.clone();
let peers = self.peers.clone();
let seq = self.seq.clone();
let config = config.clone();
tokio::spawn(async move {
handlers::heartbeat_loop(socket, peers, seq, config).await;
});
}
{
let socket = socket.clone();
let peers = self.peers.clone();
let seq = self.seq.clone();
let state = self.state_manager.clone();
let config = config.clone();
let start_time = self.start_time;
tokio::spawn(async move {
handlers::publish_status_loop(socket, peers, seq, state, config, start_time).await;
});
}
{
let socket = socket.clone();
let peers = self.peers.clone();
let seq = self.seq.clone();
let state = self.state_manager.clone();
let config = config.clone();
let gimbal = self.gimbal.clone();
let start_time = self.start_time;
let notify_lost = self.notify_lost.clone();
tokio::spawn(async move {
handlers::attitude_loop(
socket,
peers,
seq,
state,
config,
gimbal,
start_time,
notify_lost,
)
.await;
});
}
let mut buf = [0u8; 1024];
let gimbal = self.gimbal.clone();
let arbitrator = self.arbitrator.clone();
let state = self.state_manager.clone();
let peers = self.peers.clone();
let seq = self.seq.clone();
let mavlink_config = config.clone();
let start_time = self.start_time;
loop {
match socket.recv_from(&mut buf).await {
Ok((len, src_addr)) => {
debug!("Received {} bytes from {}", len, src_addr);
peers.lock().await.insert(src_addr);
let mut reader = mavlink::peek_reader::PeekReader::new(&buf[..len]);
match mavlink::read_v2_msg(&mut reader) {
Ok((header, msg)) => {
debug!("Received MAVLink message: {:?}", msg);
let ctx = RxCtx {
socket: &socket,
config: &mavlink_config,
seq: &seq,
state: &state,
arbitrator: &arbitrator,
gimbal: &gimbal,
start_time,
};
handlers::handle_mavlink_message(
header.system_id,
header.component_id,
src_addr,
msg,
&ctx,
)
.await;
}
Err(e) => debug!("Failed to parse MAVLink message: {}", e),
}
}
Err(e) => error!("UDP receive error: {}", e),
}
}
}
}