feat: HELIOS Remote v5.0.0 Final Release & Systemd Auto-Start Daemon
This commit is contained in:
13
relay/Cargo.toml
Normal file
13
relay/Cargo.toml
Normal file
@@ -0,0 +1,13 @@
|
||||
[package]
|
||||
name = "relay"
|
||||
version = "0.1.0"
|
||||
edition = "2021"
|
||||
|
||||
[dependencies]
|
||||
common = { path = "../common" }
|
||||
tokio = { version = "1.0", features = ["full"] }
|
||||
tracing = "0.1"
|
||||
tracing-subscriber = "0.3"
|
||||
serde_json = "1.0"
|
||||
uuid = { version = "1.0", features = ["v4", "serde"] }
|
||||
|
||||
242
relay/src/main.rs
Normal file
242
relay/src/main.rs
Normal file
@@ -0,0 +1,242 @@
|
||||
use common::ProtocolMessage;
|
||||
use std::collections::HashMap;
|
||||
use std::net::SocketAddr;
|
||||
use std::sync::Arc;
|
||||
use tokio::io::{AsyncBufReadExt, AsyncWriteExt, BufReader};
|
||||
use tokio::net::{TcpListener, TcpStream};
|
||||
use tokio::sync::{mpsc, Mutex};
|
||||
use tracing::{error, info, warn};
|
||||
use uuid::Uuid;
|
||||
|
||||
type MessageSender = mpsc::UnboundedSender<ProtocolMessage>;
|
||||
|
||||
#[derive(Clone, Default)]
|
||||
struct SessionManager {
|
||||
// Registered active agents: device_id -> sender
|
||||
agents: Arc<Mutex<HashMap<Uuid, MessageSender>>>,
|
||||
// Registered active viewers: device_id -> list of viewer senders
|
||||
viewers: Arc<Mutex<HashMap<Uuid, Vec<MessageSender>>>>,
|
||||
}
|
||||
|
||||
impl SessionManager {
|
||||
async fn register_agent(&self, device_id: Uuid, tx: MessageSender) {
|
||||
self.agents.lock().await.insert(device_id, tx);
|
||||
info!("Registered Agent session for device: {}", device_id);
|
||||
}
|
||||
|
||||
async fn unregister_agent(&self, device_id: &Uuid) {
|
||||
self.agents.lock().await.remove(device_id);
|
||||
info!("Unregistered Agent session for device: {}", device_id);
|
||||
}
|
||||
|
||||
async fn register_viewer(&self, device_id: Uuid, tx: MessageSender) {
|
||||
let mut viewers = self.viewers.lock().await;
|
||||
viewers.entry(device_id).or_default().push(tx);
|
||||
info!("Registered Viewer session for device: {}", device_id);
|
||||
}
|
||||
|
||||
async fn forward_to_viewers(&self, device_id: &Uuid, msg: &ProtocolMessage) {
|
||||
let viewers = self.viewers.lock().await;
|
||||
if let Some(list) = viewers.get(device_id) {
|
||||
for tx in list {
|
||||
let _ = tx.send(msg.clone());
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
async fn forward_to_agent(&self, device_id: &Uuid, msg: &ProtocolMessage) {
|
||||
let agents = self.agents.lock().await;
|
||||
if let Some(tx) = agents.get(device_id) {
|
||||
let _ = tx.send(msg.clone());
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
#[tokio::main]
|
||||
async fn main() -> Result<(), Box<dyn std::error::Error>> {
|
||||
tracing_subscriber::fmt::init();
|
||||
|
||||
let addr = "0.0.0.0:4000";
|
||||
let listener = TcpListener::bind(addr).await?;
|
||||
info!("HELIOS Relay Server listening on {}", addr);
|
||||
|
||||
let session_manager = SessionManager::default();
|
||||
|
||||
loop {
|
||||
let (socket, peer_addr) = listener.accept().await?;
|
||||
info!("New TCP connection from {}", peer_addr);
|
||||
|
||||
let sessions = session_manager.clone();
|
||||
tokio::spawn(async move {
|
||||
if let Err(e) = handle_connection(socket, peer_addr, sessions).await {
|
||||
error!("Connection error with {}: {}", peer_addr, e);
|
||||
}
|
||||
});
|
||||
}
|
||||
}
|
||||
|
||||
async fn handle_connection(
|
||||
stream: TcpStream,
|
||||
peer_addr: SocketAddr,
|
||||
sessions: SessionManager,
|
||||
) -> Result<(), Box<dyn std::error::Error>> {
|
||||
let (reader, mut writer) = stream.into_split();
|
||||
let mut lines = BufReader::new(reader).lines();
|
||||
|
||||
// 1. Read Handshake line
|
||||
let first_line = match lines.next_line().await? {
|
||||
Some(line) => line,
|
||||
None => return Ok(()),
|
||||
};
|
||||
|
||||
let handshake_msg: ProtocolMessage = serde_json::from_str(&first_line)?;
|
||||
|
||||
let (device_id, is_agent) = match handshake_msg {
|
||||
ProtocolMessage::Handshake {
|
||||
device_id,
|
||||
session_token,
|
||||
is_agent,
|
||||
} => {
|
||||
if session_token.trim().is_empty() {
|
||||
warn!("Rejected connection from {}: Empty session token", peer_addr);
|
||||
let response = ProtocolMessage::AuthResponse {
|
||||
success: false,
|
||||
message: "Invalid or missing session token".to_string(),
|
||||
};
|
||||
let mut resp_str = serde_json::to_string(&response)?;
|
||||
resp_str.push('\n');
|
||||
writer.write_all(resp_str.as_bytes()).await?;
|
||||
return Ok(());
|
||||
}
|
||||
(device_id, is_agent)
|
||||
}
|
||||
_ => {
|
||||
warn!("First packet was not Handshake from {}", peer_addr);
|
||||
return Ok(());
|
||||
}
|
||||
};
|
||||
|
||||
|
||||
// Send Handshake response
|
||||
let response = ProtocolMessage::AuthResponse {
|
||||
success: true,
|
||||
message: format!(
|
||||
"Successfully registered as {}",
|
||||
if is_agent { "Agent" } else { "Viewer" }
|
||||
),
|
||||
};
|
||||
let mut resp_str = serde_json::to_string(&response)?;
|
||||
resp_str.push('\n');
|
||||
writer.write_all(resp_str.as_bytes()).await?;
|
||||
|
||||
// Create channel for outgoing messages to this client
|
||||
let (tx, mut rx) = mpsc::unbounded_channel::<ProtocolMessage>();
|
||||
|
||||
if is_agent {
|
||||
sessions.register_agent(device_id, tx).await;
|
||||
} else {
|
||||
sessions.register_viewer(device_id, tx).await;
|
||||
}
|
||||
|
||||
// Spawn writer task
|
||||
let writer_handle = tokio::spawn(async move {
|
||||
while let Some(msg) = rx.recv().await {
|
||||
if let Ok(mut json) = serde_json::to_string(&msg) {
|
||||
json.push('\n');
|
||||
if writer.write_all(json.as_bytes()).await.is_err() {
|
||||
break;
|
||||
}
|
||||
}
|
||||
}
|
||||
});
|
||||
|
||||
// Reader loop
|
||||
while let Ok(Some(line)) = lines.next_line().await {
|
||||
if line.trim().is_empty() {
|
||||
continue;
|
||||
}
|
||||
|
||||
if let Ok(msg) = serde_json::from_str::<ProtocolMessage>(&line) {
|
||||
match &msg {
|
||||
ProtocolMessage::ScreenFrame { .. } => {
|
||||
if is_agent {
|
||||
sessions.forward_to_viewers(&device_id, &msg).await;
|
||||
}
|
||||
}
|
||||
ProtocolMessage::InputEvent { .. } => {
|
||||
if !is_agent {
|
||||
sessions.forward_to_agent(&device_id, &msg).await;
|
||||
}
|
||||
}
|
||||
ProtocolMessage::FileChunk { .. } | ProtocolMessage::ClipboardSync { .. } => {
|
||||
if is_agent {
|
||||
sessions.forward_to_viewers(&device_id, &msg).await;
|
||||
} else {
|
||||
sessions.forward_to_agent(&device_id, &msg).await;
|
||||
}
|
||||
}
|
||||
ProtocolMessage::AudioChunk { .. } | ProtocolMessage::MonitorInfoList { .. } | ProtocolMessage::TerminalOutput { .. } | ProtocolMessage::SystemMetrics { .. } | ProtocolMessage::McpResponse { .. } | ProtocolMessage::AiReport(_) | ProtocolMessage::AuditEvent(_) => {
|
||||
if is_agent {
|
||||
sessions.forward_to_viewers(&device_id, &msg).await;
|
||||
}
|
||||
}
|
||||
|
||||
ProtocolMessage::SelectMonitor { .. } | ProtocolMessage::TerminalCommand { .. } | ProtocolMessage::McpRequest { .. } => {
|
||||
if !is_agent {
|
||||
sessions.forward_to_agent(&device_id, &msg).await;
|
||||
}
|
||||
}
|
||||
ProtocolMessage::WolRequest { mac_address, broadcast_ip } => {
|
||||
|
||||
info!("Relay received WolRequest for MAC: {}", mac_address);
|
||||
let target_ip = broadcast_ip.as_deref().unwrap_or("255.255.255.255");
|
||||
if let Err(e) = send_wol_packet(&mac_address, target_ip).await {
|
||||
error!("Failed to send WOL packet: {}", e);
|
||||
}
|
||||
}
|
||||
ProtocolMessage::Heartbeat { timestamp } => {
|
||||
info!("Relay received Heartbeat (ts: {}) from {}", timestamp, peer_addr);
|
||||
}
|
||||
|
||||
|
||||
|
||||
_ => {}
|
||||
}
|
||||
}
|
||||
}
|
||||
|
||||
if is_agent {
|
||||
sessions.unregister_agent(&device_id).await;
|
||||
}
|
||||
|
||||
writer_handle.abort();
|
||||
info!("Connection closed for {}", peer_addr);
|
||||
Ok(())
|
||||
}
|
||||
|
||||
async fn send_wol_packet(mac_str: &str, broadcast_ip: &str) -> Result<(), Box<dyn std::error::Error>> {
|
||||
let mac_clean = mac_str.replace([':', '-'], "");
|
||||
if mac_clean.len() != 12 {
|
||||
return Err("Invalid MAC address length".into());
|
||||
}
|
||||
|
||||
let mut mac_bytes = [0u8; 6];
|
||||
for i in 0..6 {
|
||||
mac_bytes[i] = u8::from_str_radix(&mac_clean[i * 2..i * 2 + 2], 16)?;
|
||||
}
|
||||
|
||||
// Build 102-byte Magic Packet (6x 0xFF + 16x MAC)
|
||||
let mut packet = vec![0xFFu8; 6];
|
||||
for _ in 0..16 {
|
||||
packet.extend_from_slice(&mac_bytes);
|
||||
}
|
||||
|
||||
let socket = tokio::net::UdpSocket::bind("0.0.0.0:0").await?;
|
||||
socket.set_broadcast(true)?;
|
||||
let target = format!("{}:9", broadcast_ip);
|
||||
socket.send_to(&packet, &target).await?;
|
||||
|
||||
info!("SUCCESS: Sent WOL Magic Packet to {} (MAC: {})", target, mac_str);
|
||||
Ok(())
|
||||
}
|
||||
|
||||
Reference in New Issue
Block a user