Skip to main content

mavlink_core/connection/tcp/
sync.rs

1//! TCP MAVLink connection
2
3use crate::Connectable;
4use crate::MAVLinkMessageRaw;
5use crate::connection::get_socket_addr;
6use crate::connection::{Connection, MavConnection};
7use crate::connection_shared::{
8    ConnectionState, next_send_header, read_message, read_raw_message, write_message,
9    write_raw_message,
10};
11use crate::peek_reader::PeekReader;
12use crate::{MavHeader, MavlinkVersion, Message};
13use core::ops::DerefMut;
14use std::io;
15use std::net::ToSocketAddrs;
16use std::net::{TcpListener, TcpStream};
17use std::sync::Mutex;
18use std::time::Duration;
19
20#[cfg(feature = "mav2-message-signing")]
21use crate::SigningConfig;
22
23use super::config::{TcpConfig, TcpMode};
24
25pub fn tcpout<T: ToSocketAddrs>(address: T) -> io::Result<TcpConnection> {
26    let addr = get_socket_addr(&address)?;
27
28    let socket = TcpStream::connect(addr)?;
29    socket.set_read_timeout(Some(Duration::from_millis(100)))?;
30
31    connection_from_stream(socket)
32}
33
34fn connection_from_stream(socket: TcpStream) -> io::Result<TcpConnection> {
35    Ok(TcpConnection {
36        reader: Mutex::new(PeekReader::new(socket.try_clone()?)),
37        writer: Mutex::new(TcpWrite {
38            socket,
39            sequence: 0,
40        }),
41        state: ConnectionState::new(),
42    })
43}
44
45pub fn tcpin<T: ToSocketAddrs>(address: T) -> io::Result<TcpConnection> {
46    let addr = get_socket_addr(&address)?;
47    let listener = TcpListener::bind(addr)?;
48
49    //For now we only accept one incoming stream: this blocks until we get one
50    for incoming in listener.incoming() {
51        match incoming {
52            Ok(socket) => {
53                return Ok(TcpConnection {
54                    reader: Mutex::new(PeekReader::new(socket.try_clone()?)),
55                    writer: Mutex::new(TcpWrite {
56                        socket,
57                        sequence: 0,
58                    }),
59                    state: ConnectionState::new(),
60                });
61            }
62            Err(e) => {
63                //TODO don't println in lib
64                println!("listener err: {e}");
65            }
66        }
67    }
68    Err(io::Error::new(
69        io::ErrorKind::NotConnected,
70        "No incoming connections!",
71    ))
72}
73
74fn accept(listener: TcpListener) -> io::Result<TcpConnection> {
75    for incoming in listener.incoming() {
76        match incoming {
77            Ok(socket) => return connection_from_stream(socket),
78            Err(e) => println!("listener err: {e}"),
79        }
80    }
81    Err(io::Error::new(
82        io::ErrorKind::NotConnected,
83        "No incoming connections!",
84    ))
85}
86
87pub struct TcpConnection {
88    reader: Mutex<PeekReader<TcpStream>>,
89    writer: Mutex<TcpWrite>,
90    state: ConnectionState,
91}
92
93struct TcpWrite {
94    socket: TcpStream,
95    sequence: u8,
96}
97
98impl<M: Message> MavConnection<M> for TcpConnection {
99    fn recv(&self) -> Result<(MavHeader, M), crate::error::MessageReadError> {
100        let mut reader = self.reader.lock().unwrap();
101        read_message::<M, _>(reader.deref_mut(), &self.state)
102    }
103
104    fn recv_raw(&self) -> Result<MAVLinkMessageRaw, crate::error::MessageReadError> {
105        let mut reader = self.reader.lock().unwrap();
106        read_raw_message::<M, _>(reader.deref_mut(), &self.state)
107    }
108
109    fn try_recv(&self) -> Result<(MavHeader, M), crate::error::MessageReadError> {
110        let mut reader = self.reader.lock().unwrap();
111        reader.reader_mut().set_nonblocking(true)?;
112
113        let result = read_message::<M, _>(reader.deref_mut(), &self.state);
114
115        reader.reader_mut().set_nonblocking(false)?;
116
117        result
118    }
119
120    fn send(&self, header: &MavHeader, data: &M) -> Result<usize, crate::error::MessageWriteError> {
121        let mut lock = self.writer.lock().unwrap();
122
123        let header = next_send_header(&mut lock.sequence, header);
124        write_message(&mut lock.socket, &self.state, header, data)
125    }
126
127    fn send_raw(&self, data: &MAVLinkMessageRaw) -> Result<usize, crate::error::MessageWriteError> {
128        let mut lock = self.writer.lock().unwrap();
129        write_raw_message(&mut lock.socket, data)
130    }
131
132    fn set_protocol_version(&mut self, version: MavlinkVersion) {
133        self.state.set_protocol_version(version);
134    }
135
136    fn protocol_version(&self) -> MavlinkVersion {
137        self.state.protocol_version()
138    }
139
140    fn set_allow_recv_any_version(&mut self, allow: bool) {
141        self.state.set_allow_recv_any_version(allow);
142    }
143
144    fn allow_recv_any_version(&self) -> bool {
145        self.state.allow_recv_any_version()
146    }
147
148    #[cfg(feature = "mav2-message-signing")]
149    fn setup_signing(&mut self, signing_data: Option<SigningConfig>) {
150        self.state.setup_signing(signing_data);
151    }
152}
153
154impl Connectable for TcpConfig {
155    fn connect<M: Message>(&self) -> io::Result<Connection<M>> {
156        let conn = match self.mode {
157            TcpMode::TcpIn => match self.take_listener()? {
158                Some(listener) => accept(listener),
159                None => tcpin(&self.address),
160            },
161            TcpMode::TcpOut => match self.take_stream()? {
162                Some(stream) => connection_from_stream(stream),
163                None => tcpout(&self.address),
164            },
165        };
166
167        Ok(conn?.into())
168    }
169}