dora-cli 1.0.0-rc.5

`dora` goal is to be a low latency, composable, and distributed data flow.
use std::io::Write;

use clap::Args;
use serde::Serialize;
use tabwriter::TabWriter;
use uuid::Uuid;

use crate::{
    command::{Executable, default_tracing},
    common::{CoordinatorOptions, expect_reply, send_control_request},
    formatting::OutputFormat,
    ws_client::WsSession,
};
use dora_message::{cli_to_coordinator::ControlRequest, coordinator_to_cli::NodeInfo};

/// List all currently running nodes and their status.
///
/// Examples:
///
/// List all nodes:
///   dora node list
///
/// List nodes in a specific dataflow:
///   dora node list --dataflow my-dataflow
///
/// List nodes as JSON:
///   dora node list --format json
#[derive(Debug, Args)]
#[clap(verbatim_doc_comment)]
pub struct List {
    /// Filter by dataflow name or UUID
    #[clap(long, short = 'd', value_name = "NAME_OR_UUID")]
    pub dataflow: Option<String>,

    /// Output format
    ///
    /// `json` emits JSON Lines (one object per line). All fields are
    /// formatted strings; the `dataflow` field is omitted from each
    /// object when `--dataflow` is given.
    #[clap(long, short = 'f', value_name = "FORMAT", default_value_t = OutputFormat::Table)]
    pub format: OutputFormat,

    /// Only print node IDs, one per line
    #[clap(long, short = 'q', conflicts_with = "format")]
    pub quiet: bool,

    #[clap(flatten)]
    coordinator: CoordinatorOptions,
}

impl Executable for List {
    fn execute(self) -> eyre::Result<()> {
        default_tracing()?;

        let session = self.coordinator.connect()?;
        list(&session, self.dataflow, self.format, self.quiet)
    }
}

#[derive(Serialize)]
struct OutputEntry {
    node: String,
    status: String,
    pid: String,
    cpu: String,
    memory: String,
    restarts: String,
    #[serde(skip_serializing_if = "Option::is_none")]
    dataflow: Option<String>,
}

fn list(
    session: &WsSession,
    dataflow_filter: Option<String>,
    format: OutputFormat,
    quiet: bool,
) -> eyre::Result<()> {
    // Request node information from coordinator
    let reply = send_control_request(session, &ControlRequest::GetNodeInfo)?;
    let node_infos = expect_reply!(reply, NodeInfoList(infos))?;

    // Filter by dataflow if specified
    let filtered_nodes: Vec<NodeInfo> = if let Some(ref filter) = dataflow_filter {
        // Try to parse as UUID first
        let filter_uuid = Uuid::parse_str(filter).ok();

        node_infos
            .into_iter()
            .filter(|node| {
                // Match by UUID or name
                if let Some(uuid) = filter_uuid {
                    node.dataflow_id == uuid
                } else {
                    node.dataflow_name.as_deref() == Some(filter)
                }
            })
            .collect()
    } else {
        node_infos
    };

    // Convert to output entries
    let entries: Vec<OutputEntry> = filtered_nodes
        .into_iter()
        .map(|node| {
            let (status, pid, cpu, memory, restarts) = if let Some(metrics) = &node.metrics {
                (
                    metrics.status.to_string(),
                    metrics.pid.to_string(),
                    format!("{:.1}%", metrics.cpu_usage),
                    format!("{:.0} MB", metrics.memory_mb),
                    metrics.restart_count.to_string(),
                )
            } else {
                (
                    "Unknown".to_string(),
                    "-".to_string(),
                    "-".to_string(),
                    "-".to_string(),
                    "-".to_string(),
                )
            };

            OutputEntry {
                node: node.node_id.to_string(),
                status,
                pid,
                cpu,
                memory,
                restarts,
                dataflow: if dataflow_filter.is_none() {
                    Some(
                        node.dataflow_name
                            .unwrap_or_else(|| node.dataflow_id.to_string()),
                    )
                } else {
                    None
                },
            }
        })
        .collect();

    if quiet {
        for entry in &entries {
            println!("{}", entry.node);
        }
        return Ok(());
    }

    match format {
        OutputFormat::Table => {
            let mut tw = TabWriter::new(std::io::stdout().lock());

            // Write header
            if dataflow_filter.is_none() {
                tw.write_all(b"NODE\tSTATUS\tPID\tCPU\tMEMORY\tRESTARTS\tDATAFLOW\n")?;
            } else {
                tw.write_all(b"NODE\tSTATUS\tPID\tCPU\tMEMORY\tRESTARTS\n")?;
            }

            // Write entries
            for entry in entries {
                if let Some(ref dataflow) = entry.dataflow {
                    tw.write_all(
                        format!(
                            "{}\t{}\t{}\t{}\t{}\t{}\t{}\n",
                            entry.node,
                            entry.status,
                            entry.pid,
                            entry.cpu,
                            entry.memory,
                            entry.restarts,
                            dataflow
                        )
                        .as_bytes(),
                    )?;
                } else {
                    tw.write_all(
                        format!(
                            "{}\t{}\t{}\t{}\t{}\t{}\n",
                            entry.node,
                            entry.status,
                            entry.pid,
                            entry.cpu,
                            entry.memory,
                            entry.restarts
                        )
                        .as_bytes(),
                    )?;
                }
            }
            tw.flush()?;
        }
        OutputFormat::Json => {
            for entry in entries {
                println!("{}", serde_json::to_string(&entry)?);
            }
        }
    }

    Ok(())
}

#[cfg(test)]
mod tests {
    use super::OutputEntry;

    fn entry(node: &str, dataflow: Option<&str>) -> OutputEntry {
        OutputEntry {
            node: node.into(),
            status: "Running".into(),
            pid: "123".into(),
            cpu: "0.1%".into(),
            memory: "10 MB".into(),
            restarts: "0".into(),
            dataflow: dataflow.map(Into::into),
        }
    }

    /// Pins the runtime JSON contract the help text promises: each entry
    /// serializes to a single line that parses as an independent JSON
    /// object, and the `dataflow` key is omitted (not null) when the
    /// listing is filtered to one dataflow.
    #[test]
    fn json_entries_serialize_as_independent_single_line_objects() {
        let unfiltered = serde_json::to_string(&entry("a", Some("df"))).expect("serialize entry");
        let filtered = serde_json::to_string(&entry("b", None)).expect("serialize entry");

        for line in [&unfiltered, &filtered] {
            assert!(!line.contains('\n'), "JSONL entries must be single-line");
            serde_json::from_str::<serde_json::Value>(line)
                .expect("each output line must parse as standalone JSON");
        }

        let unfiltered: serde_json::Value =
            serde_json::from_str(&unfiltered).expect("parse unfiltered entry");
        assert_eq!(unfiltered["dataflow"], "df");
        let filtered: serde_json::Value =
            serde_json::from_str(&filtered).expect("parse filtered entry");
        assert!(
            filtered.get("dataflow").is_none(),
            "`dataflow` must be omitted, not null, when filtering by dataflow"
        );
    }
}