mavlink_core/connection/tcp/
sync.rs1use 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 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 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}