use std::{env::args, path::PathBuf};
use serde::{Deserialize, Serialize};
use tokio::{
io::{AsyncReadExt, AsyncWriteExt},
net::{
UnixStream,
unix::{OwnedReadHalf, OwnedWriteHalf},
},
signal::unix::{SignalKind, signal},
};
use tracing::info;
use crate::prelude::*;
pub mod error;
pub mod prelude;
#[derive(Debug, Clone, Serialize, Deserialize)]
pub enum BuilderEvent {
Exit,
}
#[derive(Debug, Clone, Serialize, Deserialize)]
pub enum BuilderResponse {
Ack,
}
#[derive(Debug, Clone, Copy)]
pub enum Action {
Build,
Run,
}
impl TryFrom<&str> for Action {
type Error = Error;
fn try_from(value: &str) -> Result<Self> {
if value == "build" {
return Ok(Action::Build);
}
if value == "run" {
return Ok(Action::Run);
}
Err(Error::InvalidAction(String::from(value)))
}
}
impl From<Action> for &'static str {
fn from(value: Action) -> Self {
match value {
Action::Build => "build",
Action::Run => "run",
}
}
}
impl From<Action> for String {
fn from(value: Action) -> Self {
let value: &str = value.into();
Self::from(value)
}
}
#[derive(Debug, Clone)]
pub struct BuilderSdk {
board_name: String,
board_config_name: String,
config_path: String,
action: Action,
}
impl BuilderSdk {
pub async fn init<F, Fut>(event_callback: F) -> Result<Self>
where
F: Fn(Self, BuilderEvent) -> Fut + Send + Sync + 'static,
Fut: Future<Output = Result<()>> + Send + 'static,
{
let args: Vec<String> = std::env::args().into_iter().collect();
if args.len() < 6 {
return Err(Error::MissingArgs(6, args.len()));
}
let action: Action = TryFrom::<&str>::try_from(&args[1])?;
let stream = UnixStream::connect(&args[5]).await?;
let sdk = Self {
config_path: args[2].clone(),
board_name: args[3].clone(),
board_config_name: args[4].clone(),
action,
};
let sdk_loop = sdk.clone();
let mut sigint = signal(SignalKind::interrupt())?;
tokio::spawn(async move {
while sigint.recv().await.is_some() {
info!("SIGINT received");
}
});
tokio::spawn(async move { sdk_loop.start_event_loop(stream, event_callback).await });
Ok(sdk)
}
pub fn action(&self) -> Action {
self.action
}
pub fn config_path(&self) -> PathBuf {
PathBuf::from(&self.config_path)
}
pub fn board_name(&self) -> &str {
&self.board_name
}
pub fn board_config_name(&self) -> &str {
&self.board_config_name
}
fn parse_event(payload: &str) -> Result<BuilderEvent> {
Ok(serde_json::from_str(payload)?)
}
async fn start_event_loop<F, Fut>(self, stream: UnixStream, cb: F) -> Result<()>
where
F: Fn(Self, BuilderEvent) -> Fut + Send + Sync + 'static,
Fut: Future<Output = Result<()>> + Send + 'static,
{
let mut payload = String::new();
let (mut rx, mut tx) = stream.into_split();
loop {
tokio::select! {
read_result = rx.read_to_string(&mut payload) => {
match read_result {
Ok(0) => break,
Ok(n) => {
let event = BuilderSdk::parse_event(&payload)?;
info!("Received event from builder {:?}", event);
cb(self.clone(), event).await;
info!("Acking event to builder");
let response = serde_json::to_string(&BuilderResponse::Ack)?;
tx.write_all(response.as_bytes()).await;
tx.write_all(b"\n").await;
tx.flush().await;
}
Err(e) => return Err(Error::from(e)),
}
}
_ = tokio::signal::ctrl_c() => {
info!("Received Ctrl+C, shutting down...");
cb(self.clone(), BuilderEvent::Exit).await; break;
}
}
}
Ok(())
}
}