use clap::Args;
use crate::{
command::{Executable, default_tracing, topic::selector::DataflowSelector},
common::CoordinatorOptions,
ws_client::WsSession,
};
use dora_message::{
cli_to_coordinator::ControlRequest,
coordinator_to_cli::ControlRequestReply,
id::{DataId, NodeId},
};
use eyre::{Context, bail};
#[derive(Debug, Args)]
#[clap(verbatim_doc_comment)]
pub struct Pub {
#[clap(flatten)]
selector: DataflowSelector,
#[clap(value_name = "TOPIC")]
topic: String,
#[clap(value_name = "DATA", required_unless_present = "file")]
data: Option<String>,
#[clap(long, conflicts_with = "data")]
file: Option<String>,
#[clap(long, default_value_t = 1, value_parser = clap::value_parser!(u32).range(1..))]
count: u32,
#[clap(flatten)]
coordinator: CoordinatorOptions,
}
impl Executable for Pub {
fn execute(self) -> eyre::Result<()> {
default_tracing()?;
let session = self.coordinator.connect()?;
let (dataflow_id, _descriptor) = self.selector.resolve(&session)?;
let (node_id, output_id) = parse_topic(&self.topic)?;
let data_json = if let Some(data) = self.data {
data
} else if let Some(path) = self.file {
std::fs::read_to_string(&path)
.with_context(|| format!("failed to read file: {path}"))?
} else {
bail!("either DATA or --file must be provided");
};
serde_json::from_str::<serde_json::Value>(&data_json).wrap_err(
"invalid JSON data\n\n \
hint: data must be valid JSON, e.g. '{\"key\": 42}' or '\"hello\"'",
)?;
for i in 0..self.count {
publish(&session, dataflow_id, &node_id, &output_id, &data_json)?;
if self.count > 1 {
eprintln!("Published message {}/{}", i + 1, self.count);
}
}
if self.count == 1 {
eprintln!("Published to {}", self.topic);
}
Ok(())
}
}
fn parse_topic(topic: &str) -> eyre::Result<(NodeId, DataId)> {
let parts: Vec<&str> = topic.splitn(2, '/').collect();
if parts.len() != 2 || parts[0].is_empty() || parts[1].is_empty() {
bail!("invalid topic format: expected 'node_id/output_id', got '{topic}'");
}
let node_id = parts[0]
.parse::<NodeId>()
.map_err(|e| eyre::eyre!("invalid node ID: {e}"))?;
let data_id = parts[1]
.parse::<DataId>()
.map_err(|e| eyre::eyre!("invalid output ID: {e}"))?;
Ok((node_id, data_id))
}
fn publish(
session: &WsSession,
dataflow_id: uuid::Uuid,
node_id: &NodeId,
output_id: &DataId,
data_json: &str,
) -> eyre::Result<()> {
let reply_raw = session
.request(
&serde_json::to_vec(&ControlRequest::TopicPublish {
dataflow_id,
node_id: node_id.clone(),
output_id: output_id.clone(),
data_json: data_json.to_string(),
})
.unwrap(),
)
.wrap_err("failed to send TopicPublish request")?;
let reply: ControlRequestReply =
serde_json::from_slice(&reply_raw).wrap_err("failed to parse reply")?;
match reply {
ControlRequestReply::TopicPublished => Ok(()),
ControlRequestReply::Error(err) => bail!("{err}"),
other => bail!("unexpected reply: {other:?}"),
}
}