#[cfg_attr(not(any(target_arch = "wasm32", test)), allow(dead_code))]
mod reconnect;
#[cfg(target_arch = "wasm32")]
mod wasm32;
use crate::commands::COMMANDS;
use exfiltrate_internal::rpc::{CommandInvocation, CommandResponse};
use std::sync::LazyLock;
#[cfg(not(target_arch = "wasm32"))]
use exfiltrate_internal::rpc::RPC;
#[cfg(not(target_arch = "wasm32"))]
use exfiltrate_internal::wire::{
ADDR, BACKOFF_DURATION, InFlightMessage, MAX_ATTACHMENTS, send_socket_frame, send_socket_rpc,
};
#[cfg(not(target_arch = "wasm32"))]
use std::net::{TcpListener, TcpStream};
pub struct Server {}
#[cfg(not(target_arch = "wasm32"))]
fn do_stream(mut stream: TcpStream) {
std::thread::Builder::new()
.name("exfiltrate::server do_stream".to_string())
.spawn(move || {
let mut in_flight_message = InFlightMessage::new();
loop {
let msg = in_flight_message.read_stream(&mut stream);
match msg {
Err(e) => {
eprintln!("Error reading inflight message: {:?}", e);
return;
}
Ok(exfiltrate_internal::wire::ReadStatus::WouldBlock) => {
std::thread::sleep(BACKOFF_DURATION);
}
Ok(exfiltrate_internal::wire::ReadStatus::Progress) => {
continue;
}
Ok(exfiltrate_internal::wire::ReadStatus::Completed(pop)) => {
let rpc = match rmp_serde::from_slice::<RPC>(&pop) {
Ok(rpc) => rpc,
Err(e) => {
eprintln!("Error parsing RPC message: {:?}", e);
return;
}
};
match rpc {
RPC::Command(command) => {
let mut response = do_command(command);
let reply_id = response.reply_id;
if response.response.attachment_count() > MAX_ATTACHMENTS {
response.success = false;
response.response = format!(
"response exceeds the {MAX_ATTACHMENTS}-attachment limit"
)
.into();
}
let attachments = response.response.split_data();
response.num_attachments = match u32::try_from(attachments.len()) {
Ok(count) => count,
Err(_) => {
eprintln!(
"Command {reply_id} produced too many attachments"
);
return;
}
};
if let Err(error) =
send_socket_rpc(RPC::CommandResponse(response), &mut stream)
{
eprintln!(
"Error replying to command {reply_id}: {error}"
);
return;
}
for attachment in attachments {
if let Err(error) = send_socket_frame(&attachment, &mut stream) {
eprintln!(
"Error sending attachment for command {reply_id}: {error}"
);
return;
}
}
}
RPC::CommandResponse(_response) => {
eprintln!("Received a command response on the server connection");
return;
}
_ => {
eprintln!("Unknown RPC variant received");
}
}
}
Ok(_) => {
eprintln!("Unknown ReadStatus variant received");
}
}
}
})
.unwrap();
}
fn do_command(command: CommandInvocation) -> CommandResponse {
for matcher in COMMANDS.lock_sync_read().iter() {
if matcher.name() == command.name {
let r = matcher.execute(command.args);
match r {
Ok(response) => return CommandResponse::new(true, response, command.reply_id),
Err(response) => return CommandResponse::new(false, response, command.reply_id),
}
}
}
let err_msg = format!("command not found: {}", command.name);
CommandResponse::new(false, err_msg.into(), command.reply_id)
}
pub static SERVER: LazyLock<Server> = LazyLock::new(Server::new);
impl Server {
fn new() -> Server {
#[cfg(not(target_arch = "wasm32"))]
{
Self::new_tcp()
}
#[cfg(target_arch = "wasm32")]
{
Self::new_web()
}
}
#[cfg(not(target_arch = "wasm32"))]
fn new_tcp() -> Server {
let listener = match TcpListener::bind(ADDR) {
Ok(listener) => listener,
Err(e) if e.kind() == std::io::ErrorKind::PermissionDenied => {
panic!(
"Permission denied to open the exfiltrate server socket. You may be running in a sandbox."
)
}
Err(e) => {
panic!("Can't open socket: {:?}", e);
}
};
eprintln!("Listening on {}", ADDR);
std::thread::Builder::new()
.name("exfiltrate::listen".to_string())
.spawn(move || {
for stream in listener.incoming() {
match stream {
Ok(stream) => {
do_stream(stream);
}
Err(e) => {
panic!("{}", e);
}
}
}
})
.unwrap();
Server {}
}
#[cfg(target_arch = "wasm32")]
fn new_web() -> Server {
wasm32::wasm32_go();
Server {}
}
}