feat: Add hyperdeck command

This commit is contained in:
2024-05-20 23:04:47 +01:00
parent ee375e34a8
commit f9cac35a29
14 changed files with 1122 additions and 174 deletions
+7 -1
View File
@@ -6,13 +6,19 @@ edition = "2021"
default-run = "hyperdeck-monitor"
[dependencies]
axum = { version = "0.6.10", features = ["macros", "ws"] }
color-eyre = "0.6.3"
futures = { version = "0.3.25" }
futures-util = "0.3.30"
serde = { version = "1.0.199", features = ["derive"] }
serde_json = "1.0.116"
tokio = { version = "1.37.0", features = ["full"] }
tokio-stream = { version = "0.1.12", features = ["sync"] }
tokio-tungstenite = "0.21.0"
tokio-util = "0.7.10"
tokio-util = { version = "0.7.11", features = ["full"] }
tower = { version = "0.4.13", features = ["full"] }
tower-http = { version = "0.4.0", features = ["full"] }
tracing = "0.1.40"
tracing-subscriber = { version = "0.3.18", features = ["env-filter"] }
url = "2.5.0"
uuid = { version = "1.2.2", features = ["serde", "v4"] }
+78 -1
View File
@@ -1,10 +1,35 @@
import { Hyperdeck, Commands } from 'hyperdeck-connection';
import WebSocket from 'ws';
interface WrappedHyperdeck {
ip: String,
port: number,
hyperdeck: Hyperdeck
}
const hyperdecks: Map<string, WrappedHyperdeck> = new Map()
enum WebSocketMessageType {
AddHyperdeck = "add_hyperdeck"
}
type WebSocketMessage = {
type: WebSocketMessageType.AddHyperdeck,
id: string,
ip: string,
port: number
}
const wss = new WebSocket.Server({ port: 7867 });
wss.on('connection', function connection(ws) {
ws.on('message', function message(data) {
console.log('received: %s', data);
try {
const message = JSON.parse(data.toString()) as Partial<WebSocketMessage>;
handle_message(message)
} catch (_err) {
return;
}
});
ws.send(JSON.stringify({
@@ -12,3 +37,55 @@ wss.on('connection', function connection(ws) {
message: "Hello"
}));
});
function exhaustiveMatch(_never: never) {
return;
}
function handle_message(message: Partial<WebSocketMessage>) {
console.log(JSON.stringify(message));
if (message.type === undefined) return;
switch (message.type) {
case WebSocketMessageType.AddHyperdeck:
if (message.id === undefined) return;
if (message.ip === undefined) return;
if (message.port === undefined) return;
if (isNaN(message.port)) return;
if (message.port <= 0) return;
console.log("Adding hyperdeck");
const newHyperdeck = new Hyperdeck()
// hyperdecks.set(message.id, {
// ip: message.ip,
// port: message.port,
// hyperdeck: newHyperdeck
// });
newHyperdeck.on('connected', (info) => {
console.log(JSON.stringify(info))
newHyperdeck.sendCommand(new Commands.TransportInfoCommand()).then((transportInfo) => {
console.log(JSON.stringify(transportInfo))
})
})
newHyperdeck.on('notify.slot', function (state) {
console.log(JSON.stringify(state)) // catch the slot state change.
})
newHyperdeck.on('notify.transport', function (state) {
console.log(JSON.stringify(state)) // catch the transport state change.
})
newHyperdeck.on('error', (err) => {
console.log('Hyperdeck error', JSON.stringify(err))
})
newHyperdeck.connect(message.ip, message.port)
break;
default:
exhaustiveMatch(message.type)
}
}
+172
View File
@@ -0,0 +1,172 @@
use axum::extract::ws::Message;
use axum::extract::{State, WebSocketUpgrade};
use axum::response::Html;
use axum::{
body::Bytes,
extract::Path,
http::{header, HeaderValue, Method},
response::IntoResponse,
routing::get,
Router,
};
use message::{ClientRequest, HyperdeckMonitorState, ServerEvent};
use serde::{Deserialize, Serialize};
use std::{
collections::HashMap,
net::{Ipv4Addr, SocketAddr},
sync::Arc,
time::Duration,
};
use tokio::sync::{Mutex, RwLock};
use tower::ServiceBuilder;
use tower_http::timeout::TimeoutLayer;
use tower_http::ServiceBuilderExt;
use tower_http::{
cors::{Any, CorsLayer},
trace::{DefaultMakeSpan, DefaultOnResponse, TraceLayer},
LatencyUnit,
};
use tracing::info;
use uuid::Uuid;
pub mod message;
mod ws;
#[derive(Debug, Clone)]
pub struct Client {
pub sender: Option<tokio::sync::broadcast::Sender<Message>>,
}
type Clients = Arc<Mutex<HashMap<Uuid, Client>>>;
pub async fn initialize_api(
mut state_rx: tokio::sync::broadcast::Receiver<HyperdeckMonitorState>,
client_request_tx: tokio::sync::mpsc::UnboundedSender<ClientRequest>,
) {
info!("Initializing API");
let clients: Clients = Default::default();
let state = Arc::new(RwLock::new(state_rx.recv().await.unwrap()));
let state_clients = clients.clone();
let state_loop = state.clone();
tokio::spawn(async move {
loop {
if let Ok(hyperdeck_monitor_state) = state_rx.recv().await {
let mut state = state_loop.write().await;
*state = hyperdeck_monitor_state.clone();
let clients = state_clients.lock().await;
let state_json = serde_json::to_string(&ServerEvent::HyperdeckMonitorState(
hyperdeck_monitor_state.into(),
))
.unwrap();
for (_, client) in clients.iter() {
if let Some(sender) = &client.sender {
let message: Message = Message::Text(state_json.clone());
let _ = sender.send(message);
}
}
}
}
});
let app_state = AppState {
state,
client_request_tx,
clients,
port: 9681,
};
let addr = SocketAddr::from((Ipv4Addr::UNSPECIFIED, app_state.port));
info!("Listening on {}", addr);
// TODO: This could fail, need to figure out how to get a result from this
let _ = axum::Server::bind(&addr)
.serve(app(app_state).into_make_service())
.await;
}
#[derive(Clone)]
struct AppState {
state: Arc<RwLock<HyperdeckMonitorState>>,
client_request_tx: tokio::sync::mpsc::UnboundedSender<ClientRequest>,
clients: Clients,
port: u16,
}
fn app(state: AppState) -> Router {
let sensitive_headers: Arc<[_]> = vec![header::AUTHORIZATION, header::COOKIE].into();
let middleware = ServiceBuilder::new()
// Mark the `Authorization` and `Cookie` headers as sensitive so it doesn't show in logs
.sensitive_request_headers(sensitive_headers.clone())
// Add high level tracing/logging to all requests
.layer(
TraceLayer::new_for_http()
.on_body_chunk(|chunk: &Bytes, latency: Duration, _: &tracing::Span| {
tracing::trace!(size_bytes = chunk.len(), latency = ?latency, "sending body chunk")
})
.make_span_with(DefaultMakeSpan::new().include_headers(true))
.on_response(DefaultOnResponse::new().include_headers(true).latency_unit(LatencyUnit::Micros)),
)
.sensitive_response_headers(sensitive_headers)
// Set a timeout
.layer(TimeoutLayer::new(Duration::from_secs(10)))
// Box the response body so it implements `Default` which is required by axum
.map_response_body(axum::body::boxed)
// Compress responses
.compression()
// Set a `Content-Type` if there isn't one already.
.insert_response_header_if_not_present(
header::CONTENT_TYPE,
HeaderValue::from_static("application/octet-stream"),
);
let cors = CorsLayer::new()
.allow_methods(vec![
Method::GET,
Method::POST,
Method::PUT,
Method::DELETE,
Method::OPTIONS,
])
.allow_headers(Any)
.allow_origin(Any)
.allow_credentials(false);
Router::new()
.route("/", get(get_index))
.route("/ws", get(upgrade_ws))
.layer(middleware)
.layer(cors)
.with_state(state)
}
#[derive(Debug, Serialize, Deserialize)]
pub struct WebSocketUpgradeRequest {}
async fn get_index() -> Html<String> {
Html(format!("Hello!"))
}
#[axum::debug_handler]
async fn upgrade_ws(state: State<AppState>, ws: WebSocketUpgrade) -> impl IntoResponse {
info!("New client websocket connection");
let client_id = uuid::Uuid::new_v4();
state
.clients
.lock()
.await
.insert(client_id.clone(), Client { sender: None });
let client = state.clients.lock().await.get(&client_id).cloned().unwrap();
ws.on_upgrade(move |socket| {
ws::client_connection(
state.client_request_tx.clone(),
socket,
client_id,
state.state.clone(),
state.clients.clone(),
client,
)
})
}
+36
View File
@@ -0,0 +1,36 @@
use std::collections::HashMap;
use serde::{Deserialize, Serialize};
#[derive(Debug, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
#[serde(tag = "type")]
pub enum ClientRequest {
AddHyperdeck(AddHyperdeckRequest),
}
#[derive(Debug, Serialize, Deserialize)]
pub struct AddHyperdeckRequest {
pub name: String,
pub ip: String,
pub port: u16,
}
#[derive(Debug, Serialize, Deserialize)]
#[serde(rename_all = "snake_case")]
#[serde(tag = "type")]
pub enum ServerEvent {
HyperdeckMonitorState(HyperdeckMonitorState),
}
#[derive(Debug, Clone, Default, Serialize, Deserialize)]
pub struct HyperdeckMonitorState {
pub hyperdecks: HashMap<String, HyperdeckState>,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub struct HyperdeckState {
pub name: String,
pub ip: String,
pub port: u16,
}
+84
View File
@@ -0,0 +1,84 @@
use std::{future, sync::Arc};
use super::message::{ClientRequest, HyperdeckMonitorState};
use crate::api::ServerEvent;
use axum::extract::ws::{Message, WebSocket};
use futures::StreamExt;
use tokio::sync::RwLock;
use tokio_stream::wrappers::BroadcastStream;
use tracing::{debug, error, log::info};
use uuid::Uuid;
use super::{Client, Clients};
pub async fn client_connection(
client_request_tx: tokio::sync::mpsc::UnboundedSender<ClientRequest>,
ws: WebSocket,
id: Uuid,
state: Arc<RwLock<HyperdeckMonitorState>>,
clients: Clients,
mut client: Client,
) {
let (client_ws_sender, mut client_ws_rcv) = ws.split();
let (client_sender, client_rcv) = tokio::sync::broadcast::channel::<Message>(10);
let client_rcv = BroadcastStream::new(client_rcv);
tokio::task::spawn(
client_rcv
.filter(|msg| future::ready(msg.is_ok()))
.map(|msg| Ok(msg.unwrap()))
.forward(client_ws_sender),
);
let current_state = state.read().await.clone();
let state_json =
serde_json::to_string(&ServerEvent::HyperdeckMonitorState(current_state.into())).unwrap();
client_sender.send(Message::Text(state_json.clone())).ok();
client.sender = Some(client_sender);
clients.lock().await.insert(id, client);
info!("{} connected", id);
while let Some(result) = client_ws_rcv.next().await {
let msg = match result {
Ok(msg) => msg,
Err(e) => {
error!("error resolving ws message for id: {}: {}", id.clone(), e);
break;
}
};
client_msg(client_request_tx.clone(), &id, msg).await;
}
clients.lock().await.remove(&id);
info!("{} disconnected", id);
}
async fn client_msg(
client_request_tx: tokio::sync::mpsc::UnboundedSender<ClientRequest>,
id: &Uuid,
msg: Message,
) {
debug!("received message from {}: {:?}", id, msg);
let message = match msg.into_text() {
Ok(v) => v,
Err(err) => {
error!("error: {:?}", err);
return;
}
};
if message == "ping" || message == "ping\n" {
return;
}
let client_request: super::message::ClientRequest = match serde_json::from_str(&message) {
Ok(v) => v,
Err(_) => {
return;
}
};
let _ = client_request_tx.send(client_request);
}
+175 -46
View File
@@ -1,5 +1,6 @@
use std::{net::IpAddr, time::Duration};
use std::{process::Stdio, time::Duration};
use api::message::{AddHyperdeckRequest, ClientRequest, HyperdeckMonitorState, HyperdeckState};
use color_eyre::Report;
use futures_util::{
pin_mut, select,
@@ -9,9 +10,14 @@ use futures_util::{
use serde::{Deserialize, Serialize};
use tokio::net::TcpStream;
use tokio_tungstenite::{tungstenite::Message, MaybeTlsStream, WebSocketStream};
use tokio_util::sync::CancellationToken;
use tokio_util::{
codec::{FramedRead, LinesCodec},
sync::CancellationToken,
};
use tracing_subscriber::EnvFilter;
mod api;
#[tokio::main]
async fn main() {
setup_logging().expect("Failed to setup logging");
@@ -20,35 +26,169 @@ async fn main() {
let cancel = CancellationToken::new();
let node_process = run_node_process(cancel.clone()).fuse();
let (ws_message_tx, ws_message_rx) = tokio::sync::mpsc::unbounded_channel();
let (commands_tx, commands_rx) = tokio::sync::mpsc::unbounded_channel();
let (node_ws_message_tx, node_ws_message_rx) = tokio::sync::mpsc::unbounded_channel();
let (node_commands_tx, node_commands_rx) = tokio::sync::mpsc::unbounded_channel();
let state = AppState::default();
let ws_process = talk_to_node_ws(state, ws_message_tx, commands_rx, cancel.clone()).fuse();
let node_ws_communication =
talk_to_node_ws(state, node_ws_message_tx, node_commands_rx, cancel.clone()).fuse();
let (state_tx, state_rx) = tokio::sync::broadcast::channel(1);
let (client_request_tx, client_request_rx) = tokio::sync::mpsc::unbounded_channel();
let api = api::initialize_api(state_rx, client_request_tx).fuse();
let hyperdeck_monitor = run(
node_commands_tx,
node_ws_message_rx,
state_tx,
client_request_rx,
cancel.clone(),
)
.fuse();
pin_mut!(node_process);
pin_mut!(ws_process);
pin_mut!(node_ws_communication);
pin_mut!(api);
pin_mut!(hyperdeck_monitor);
select! {
_ = node_process => {},
_ = ws_process => {},
_ = node_ws_communication => {},
_ = api => {},
_ = hyperdeck_monitor => {},
_ = cancel.cancelled().fuse() => {}
};
cancel.cancel();
}
async fn run(
mut node_commands_tx: tokio::sync::mpsc::UnboundedSender<NodeWsCommand>,
mut node_ws_message_rx: tokio::sync::mpsc::UnboundedReceiver<NodeWsMessageReceived>,
mut state_tx: tokio::sync::broadcast::Sender<HyperdeckMonitorState>,
mut client_request_rx: tokio::sync::mpsc::UnboundedReceiver<ClientRequest>,
cancel: CancellationToken,
) {
let mut state = HyperdeckMonitorState::default();
let _ = state_tx.send(state.clone());
let mut ping_interval = tokio::time::interval(Duration::from_millis(500));
ping_interval.set_missed_tick_behavior(tokio::time::MissedTickBehavior::Delay);
while !cancel.is_cancelled() {
let state_modified = select! {
_ = ping_interval.tick().fuse() => {
// TODO: Ping node, check it's still alive
false
},
message_from_node = node_ws_message_rx.recv().fuse() => {
if let Some(msg) = message_from_node {
handle_message_from_node(msg, &mut node_commands_tx, &mut state).await
} else {
false
}
},
message_from_client = client_request_rx.recv().fuse() => {
if let Some(msg) = message_from_client {
handle_message_from_client(msg, &mut node_commands_tx, &mut state).await
} else {
false
}
}
};
if state_modified {
let _ = state_tx.send(state.clone());
}
}
}
async fn handle_message_from_node(
msg: NodeWsMessageReceived,
node_commands_tx: &mut tokio::sync::mpsc::UnboundedSender<NodeWsCommand>,
state: &mut HyperdeckMonitorState,
) -> bool {
match msg {
NodeWsMessageReceived::Log { message } => {
tracing::info!("[NODE] {message}");
false
}
}
}
async fn handle_message_from_client(
msg: ClientRequest,
node_commands_tx: &mut tokio::sync::mpsc::UnboundedSender<NodeWsCommand>,
state: &mut HyperdeckMonitorState,
) -> bool {
match msg {
ClientRequest::AddHyperdeck(AddHyperdeckRequest { name, ip, port }) => {
tracing::info!("Adding hyperdeck");
let id = uuid::Uuid::new_v4();
state.hyperdecks.insert(
id.to_string(),
HyperdeckState {
name,
ip: ip.clone(),
port,
},
);
let _ = node_commands_tx.send(NodeWsCommand::AddHyperdeck(AddHyperdeckCommand {
id: id.to_string(),
ip,
port,
}));
true
}
}
}
async fn run_node_process(cancel: CancellationToken) {
while !cancel.is_cancelled() {
// Back-off in case we are immediately crashing in a loop.
tokio::time::sleep(Duration::from_secs(1)).await;
let result = tokio::process::Command::new("node")
.arg("monitor/index.js")
.output()
.await;
if let Ok(output) = result {
if !output.status.success() {
let err = String::from_utf8(output.stderr).unwrap_or("Unknown".to_string());
tracing::error!("Node process exited with error: {}", err);
// Back-off in case we are immediately crashing in a loop.
tokio::time::sleep(Duration::from_secs(1)).await;
.stdin(Stdio::piped())
.stdout(Stdio::piped())
.stderr(Stdio::piped())
.spawn();
match result {
Ok(mut child_process) => {
let Some(raw_stdout) = child_process.stdout.take() else {
let _ = child_process.kill().await;
continue;
};
let Some(raw_stderr) = child_process.stderr.take() else {
let _ = child_process.kill().await;
continue;
};
let mut stdout = FramedRead::new(raw_stdout, LinesCodec::new())
.map(|data| data.expect("Could not read stdout"));
let mut stderr = FramedRead::new(raw_stderr, LinesCodec::new())
.map(|data| data.expect("Could not read stderr"));
while !cancel.is_cancelled() {
select! {
line = stdout.next().fuse() => {
if let Some(line) = line {
tracing::info!("[NODE] {line}");
}
}
line = stderr.next().fuse() => {
if let Some(line) = line {
tracing::error!("[NODE] {line}");
}
}
}
}
let _ = child_process.kill().await;
}
Err(err) => {
tracing::error!("Error running Node child process: {err}");
}
}
}
@@ -56,9 +196,15 @@ async fn run_node_process(cancel: CancellationToken) {
#[derive(Default)]
struct AppState {}
#[derive(Debug, Serialize, Deserialize)]
#[serde(tag = "type")]
enum NodeWsCommand {
#[serde(rename = "ping")]
Ping,
#[serde(rename = "add_hyperdeck")]
AddHyperdeck(AddHyperdeckCommand),
#[serde(rename = "remove_hyperdeck")]
RemoveHyperdeck(RemoveHyperdeckCommand),
}
@@ -69,16 +215,18 @@ enum NodeWsMessageReceived {
Log { message: String },
}
#[derive(Serialize, Deserialize)]
#[derive(Debug, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
struct AddHyperdeckCommand {
ip: IpAddr,
id: String,
ip: String,
port: u16,
}
#[derive(Serialize, Deserialize)]
#[derive(Debug, Serialize, Deserialize)]
#[serde(rename_all = "camelCase")]
struct RemoveHyperdeckCommand {
ip: IpAddr,
id: String,
}
async fn talk_to_node_ws(
@@ -126,28 +274,13 @@ async fn handle_outbound_messages(
mut socket_tx: SplitSink<WebSocketStream<MaybeTlsStream<TcpStream>>, Message>,
) {
while let Some(command) = commands_rx.recv().await {
match command {
NodeWsCommand::Ping => {
let _ = socket_tx
.send(tokio_tungstenite::tungstenite::Message::Ping(vec![]))
.await;
}
NodeWsCommand::AddHyperdeck(command) => {
let _ = socket_tx
.send(tokio_tungstenite::tungstenite::Message::Text(
serde_json::to_string(&command)
.expect("Could not serialize AddHyperdeck command"),
))
.await;
}
NodeWsCommand::RemoveHyperdeck(command) => {
let _ = socket_tx
.send(tokio_tungstenite::tungstenite::Message::Text(
serde_json::to_string(&command)
.expect("Could not serialize RemoveHyperdeck command"),
))
.await;
}
if let Err(err) = socket_tx
.send(tokio_tungstenite::tungstenite::Message::Text(
serde_json::to_string(&command).expect("Could not serialize command"),
))
.await
{
tracing::error!("Error sending command to Node proccess: {err}");
}
}
}
@@ -161,11 +294,7 @@ async fn handle_inbound_messages(
match message {
Ok(tokio_tungstenite::tungstenite::Message::Text(text)) => {
if let Ok(received) = serde_json::from_str::<NodeWsMessageReceived>(&text) {
match received {
NodeWsMessageReceived::Log { message } => {
tracing::info!("Message from Node process: {message}");
}
}
let _ = ws_message_tx.send(received);
}
}
Ok(tokio_tungstenite::tungstenite::Message::Pong(_)) => {}