mirror of
https://github.com/hubaldv/bioz-host-rs.git
synced 2026-07-24 00:57:44 +00:00
Add basic TCP.
This commit is contained in:
60
src/tcp.rs
Normal file
60
src/tcp.rs
Normal file
@@ -0,0 +1,60 @@
|
||||
use tokio::net::TcpListener;
|
||||
use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader, BufWriter};
|
||||
use tokio::sync::mpsc::Sender;
|
||||
|
||||
use crate::signals::StartStopSignal;
|
||||
|
||||
pub async fn tcp_command_server(tx: Sender<StartStopSignal>) {
|
||||
let listener = TcpListener::bind("127.0.0.1:12345").await.unwrap();
|
||||
println!("TCP command server listening on 12345");
|
||||
|
||||
loop {
|
||||
let (socket, addr) = listener.accept().await.unwrap();
|
||||
println!("Client connected: {:?}", addr);
|
||||
|
||||
let tx = tx.clone();
|
||||
|
||||
tokio::spawn(async move {
|
||||
let (reader, writer) = socket.into_split();
|
||||
|
||||
let mut reader = BufReader::new(reader);
|
||||
let mut writer = BufWriter::new(writer);
|
||||
let mut lines = reader.lines();
|
||||
|
||||
while let Ok(Some(line)) = lines.next_line().await {
|
||||
println!("Received: {}", line);
|
||||
|
||||
match parse_command(&line) {
|
||||
Some(signal) => {
|
||||
if tx.send(signal).await.is_ok() {
|
||||
let _ = writer.write_all(b"OK\n").await;
|
||||
} else {
|
||||
let _ = writer.write_all(b"ERR internal channel failure\n").await;
|
||||
break;
|
||||
}
|
||||
}
|
||||
None => {
|
||||
let _ = writer.write_all(b"ERR unknown command\n").await;
|
||||
}
|
||||
}
|
||||
|
||||
let _ = writer.flush().await;
|
||||
}
|
||||
|
||||
println!("Client disconnected: {:?}", addr);
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
fn parse_command(line: &str) -> Option<StartStopSignal> {
|
||||
let parts: Vec<&str> = line.split_whitespace().collect();
|
||||
|
||||
match parts.as_slice() {
|
||||
["STOP"] => Some(StartStopSignal::Stop),
|
||||
|
||||
_ => {
|
||||
eprintln!("Unknown command: {}", line);
|
||||
None
|
||||
}
|
||||
}
|
||||
}
|
||||
Reference in New Issue
Block a user