use std::path::PathBuf;
use nu_engine::CallExt;
use nu_protocol::{
engine::{Call, Command, EngineState, Stack},
shell_error::generic::GenericError,
LabeledError, PipelineData, ShellError, Signature, SyntaxShape, Type, Value,
};
use zenoh::{session, Wait};
use crate::{call_ext2::CallExt2, conv, signature_ext::SignatureExt, State};
#[derive(Clone)]
pub(crate) struct Open {
state: State,
}
impl Open {
pub(crate) fn new(state: State) -> Self {
Self { state }
}
}
impl Command for Open {
fn name(&self) -> &str {
"zenoh session open"
}
fn signature(&self) -> Signature {
let sig = Signature::build(self.name())
.session()
.zenoh_category()
.config()
.named(
"config-file",
SyntaxShape::Filepath,
"Path to a Zenoh configuration file",
None,
)
.input_output_type(Type::Nothing, Type::Nothing);
if self.state.options.experimental_options {
sig.named("runtime", SyntaxShape::String, "Runtime name", None)
} else {
sig
}
}
fn description(&self) -> &str {
"Open or re-open a session"
}
fn run(
&self,
engine_state: &EngineState,
stack: &mut Stack,
call: &Call,
_input: PipelineData,
) -> Result<PipelineData, ShellError> {
let file_path = call.get_flag::<PathBuf>(engine_state, stack, "config-file")?;
let runtime_name = call.get_flag::<String>(engine_state, stack, "runtime")?;
let config_record = call.opt::<Value>(engine_state, stack, 0)?;
let config = match (
file_path.as_ref(),
config_record.as_ref(),
runtime_name.as_ref(),
) {
(Some(file_path), None, None) => zenoh::Config::from_file(file_path).map_err(|e| {
nu_protocol::LabeledError::new("Failed to load config file").with_label(
format!("Could not read config from {}: {}", file_path.display(), e),
call.head,
)
})?,
(None, Some(config_record), None) => match config_record {
val @ Value::Record { .. } => {
let json_value =
conv::value_to_json_value(engine_state, val, call.head, false)?;
zenoh::Config::from_json5(&json_value.to_string()).map_err(|e| {
nu_protocol::LabeledError::new("Failed to parse config record")
.with_label(format!("Could not parse config record: {e}"), call.head)
})?
}
_ => {
return Err(ShellError::Generic(
GenericError::new(
"Invalid config type",
"Config must be a record",
call.head,
)
.with_help("Provide a record with Zenoh configuration options"),
));
}
},
(None, None, Some(runtime_name)) => {
let runtime = self
.state
.runtimes
.read()
.unwrap()
.get(runtime_name)
.ok_or_else(|| {
LabeledError::new(format!("runtime '{runtime_name}' was not found"))
})?
.clone();
let session_name = call.session(engine_state, stack)?;
let mut sessions = self.state.sessions.write().unwrap();
if let Some(sess) = sessions.remove(&session_name) {
sess.close().wait().map_err(|e| {
nu_protocol::LabeledError::new(
"Failed to reopen Zenoh session '{session_name}'",
)
.with_label(format!("Could not close Zenoh session: {e}"), call.head)
})?
}
let new_session = session::init(runtime.into()).wait().map_err(|e| {
nu_protocol::LabeledError::new("Failed to open Zenoh session")
.with_label(format!("Could not establish Zenoh session: {e}"), call.head)
})?;
sessions.insert(session_name, new_session);
return Ok(PipelineData::Value(Value::nothing(call.head), None));
}
(None, None, None) => zenoh::Config::default(),
_ => {
return Err(ShellError::Generic(GenericError::new(
"Conflicting arguments",
"Only one of RECORD, --config-file or --runtime can be specified",
call.head,
)));
}
};
let session_name = call.session(engine_state, stack)?;
let mut sessions = self.state.sessions.write().unwrap();
if let Some(sess) = sessions.remove(&session_name) {
sess.close().wait().map_err(|e| {
nu_protocol::LabeledError::new("Failed to reopen Zenoh session '{session_name}'")
.with_label(format!("Could not close Zenoh session: {e}"), call.head)
})?
}
let new_session = zenoh::open(config).wait().map_err(|e| {
nu_protocol::LabeledError::new("Failed to open Zenoh session")
.with_label(format!("Could not establish Zenoh session: {e}"), call.head)
})?;
sessions.insert(session_name, new_session);
Ok(PipelineData::Value(Value::nothing(call.head), None))
}
}