use crate::types::state::{StateEvent, StateQueryKind};
use crate::types::wire::{
self, Error as WireError, QueryPayload, Request, Response, ScriptStateRequest,
ScriptStateResponse, SendPort, Sender, service_name,
};
use futures_lite::StreamExt;
use crate::errors::{Error, Result};
fn connect() -> Result<Sender<Request>> {
wire::connect(&service_name()).map_err(|err| match err {
WireError::NotRunning => Error::Generic("paneru is not running".to_string()),
other => Error::from(other),
})
}
pub async fn send_command(argv: impl IntoIterator<Item = String>) -> Result<()> {
let argv = argv.into_iter().collect::<Vec<_>>();
let borrowed = argv.iter().map(String::as_str).collect::<Vec<_>>();
let command = crate::types::argv::parse_command(&borrowed)?;
connect()?.send(&Request::Command(command)).await?;
Ok(())
}
pub async fn query(kind: StateQueryKind) -> Result<String> {
let response: Response = connect()?.call(&Request::Query(kind)).await?;
match response {
Response::Query(payload) => render(&payload),
other => Err(unexpected(&other)),
}
}
pub async fn script_state(request: ScriptStateRequest) -> Result<String> {
let response: Response = connect()?.call(&Request::ScriptState(request)).await?;
let answer = match response {
Response::ScriptState(answer) => answer,
other => return Err(unexpected(&other)),
};
let value = match answer {
ScriptStateResponse::Value(value) => serde_json::json!({
"value": value.map(serde_json::Value::from),
}),
ScriptStateResponse::Write(outcome) => outcome
.to_json()
.map_err(|err| Error::Generic(err.to_string()))?,
};
Ok(value.to_string())
}
pub async fn subscribe() -> Result<()> {
use std::io::Write;
let events = connect()?
.subscribe::<StateEvent>(&Request::Subscribe)
.await?;
let mut events = std::pin::pin!(events);
while let Some(delivery) = events.next().await {
let event = match delivery {
Ok(delivery) => delivery.value,
Err(WireError::PeerGone) => break,
Err(err) => return Err(Error::from(err)),
};
let line = event
.to_json()
.map_err(|err| Error::Generic(err.to_string()))?
.to_string();
let mut stdout = std::io::stdout();
if writeln!(stdout, "{line}")
.and_then(|()| stdout.flush())
.is_err()
{
break;
}
}
Ok(())
}
pub fn run(command: ClientCommand) -> Result<()> {
futures_lite::future::block_on(async move {
match command {
ClientCommand::Send(argv) => send_command(argv).await,
ClientCommand::Query(kind) => {
println!("{}", query(kind).await?);
Ok(())
}
ClientCommand::ScriptState(request) => {
println!("{}", script_state(request).await?);
Ok(())
}
ClientCommand::Subscribe => subscribe().await,
}
})
}
#[derive(Debug)]
pub enum ClientCommand {
Send(Vec<String>),
Query(StateQueryKind),
ScriptState(ScriptStateRequest),
Subscribe,
}
fn render(payload: &QueryPayload) -> Result<String> {
payload
.to_json()
.map(|value| value.to_string())
.map_err(|err| Error::Generic(err.to_string()))
}
fn unexpected(response: &Response) -> Error {
match response {
Response::Error(message) => Error::Generic(message.clone()),
other => Error::Generic(format!("unexpected response: {other:?}")),
}
}