Skip to main content

viva_camctl/
cmd_bench.rs

1use std::fs::File;
2use std::net::Ipv4Addr;
3use std::path::PathBuf;
4use std::time::{Duration, SystemTime};
5
6use anyhow::{Context, Result, anyhow, bail};
7use bytes::BytesMut;
8use serde::Serialize;
9use tokio::time::{self, Instant, MissedTickBehavior};
10use tracing::{info, warn};
11
12use viva_genicam::gige::gvsp::{self, GvspPacket};
13use viva_genicam::pfnc::PixelFormat;
14use viva_genicam::{Frame, StreamBuilder, StreamDest, parse_chunk_bytes};
15
16use viva_gige::nic::IfaceSelector;
17
18use crate::common::{self, DEFAULT_DISCOVERY_TIMEOUT_MS};
19
20#[derive(Debug, Clone)]
21pub struct BenchArgs {
22    pub ip: Option<Ipv4Addr>,
23    pub index: Option<usize>,
24    pub iface: Option<IfaceSelector>,
25    pub mode: String,
26    pub group: Option<Ipv4Addr>,
27    pub port: u16,
28    pub duration_s: u64,
29    pub json_out: Option<PathBuf>,
30}
31
32#[derive(Debug, Serialize)]
33struct BenchReport {
34    duration_s: u64,
35    frames: u64,
36    bytes: u64,
37    avg_fps: f64,
38    avg_mbps: f64,
39    drops: u64,
40    resends: u64,
41    mode: String,
42}
43
44struct BlockState {
45    block_id: u64,
46    width: u32,
47    height: u32,
48    pixel_format: PixelFormat,
49    timestamp: u64,
50    payload: BytesMut,
51}
52
53pub async fn run(args: BenchArgs, emit_json: bool) -> Result<()> {
54    let timeout = Duration::from_millis(DEFAULT_DISCOVERY_TIMEOUT_MS);
55    let device = common::select_device(args.ip, args.index, args.iface.as_ref(), timeout).await?;
56    info!(ip = %device.ip, "opening camera for benchmark");
57    let mut camera = common::open_camera(&device)
58        .await
59        .context("open camera for bench")?;
60    let mut stream_device = common::open_control(&device)
61        .await
62        .context("open control channel for bench")?;
63
64    let iface = common::resolve_receive_iface(args.iface.as_ref(), device.ip)?;
65    let host_ip = iface
66        .ipv4()
67        .ok_or_else(|| anyhow!("interface {} has no IPv4 address", iface.name()))?;
68    let mode = parse_mode(&args.mode)?;
69
70    if let StreamMode::Multicast = mode {
71        let group = args
72            .group
73            .ok_or_else(|| anyhow!("multicast mode requires --group"))?;
74        camera
75            .configure_stream_multicast(0, group, args.port)
76            .context("configure multicast destination")?;
77    }
78
79    let mut builder = StreamBuilder::new(&mut stream_device).iface(iface.clone());
80    let dest = match mode {
81        StreamMode::Unicast => StreamDest::Unicast {
82            dst_ip: host_ip,
83            dst_port: args.port,
84        },
85        StreamMode::Multicast => {
86            let group = args
87                .group
88                .ok_or_else(|| anyhow!("multicast mode requires --group"))?;
89            StreamDest::Multicast {
90                group,
91                port: args.port,
92                loopback: false,
93                ttl: 1,
94            }
95        }
96    };
97    builder = builder.dest(dest);
98    builder = builder.auto_packet_size();
99    let stream = builder.build().await.context("negotiate stream")?;
100
101    camera.acquisition_start().context("start acquisition")?;
102    let mut recv_buffer = vec![0u8; (stream.params().packet_size as usize + 64).max(4096)];
103    let stats = stream.stats_handle();
104    let mut state: Option<BlockState> = None;
105    let mut ctrl_c = Box::pin(tokio::signal::ctrl_c());
106    let mut ticker = time::interval(Duration::from_secs(1));
107    ticker.set_missed_tick_behavior(MissedTickBehavior::Delay);
108    let duration = Duration::from_secs(args.duration_s.max(1));
109    let end_deadline = Instant::now() + duration;
110    let mut interrupted = false;
111
112    loop {
113        if Instant::now() >= end_deadline {
114            info!("benchmark duration elapsed");
115            break;
116        }
117
118        tokio::select! {
119            _ = ticker.tick() => {
120                let snapshot = stream.stats();
121                println!(
122                    "[bench] fps={:.1} Mbps={:.2} frames={} drops={} resends={}",
123                    snapshot.avg_fps,
124                    snapshot.avg_mbps,
125                    snapshot.frames,
126                    snapshot.drops,
127                    snapshot.resends,
128                );
129            }
130            _ = &mut ctrl_c => {
131                info!("received ctrl-c; stopping bench early");
132                interrupted = true;
133                break;
134            }
135            recv = stream.socket().expect("UDP socket").recv_from(&mut recv_buffer) => {
136                let (len, _) = match recv {
137                    Ok(result) => result,
138                    Err(err) => {
139                        warn!(error = %err, "socket receive failed");
140                        stats.record_drop();
141                        continue;
142                    }
143                };
144                let packet = match gvsp::parse_packet(&recv_buffer[..len]) {
145                    Ok(packet) => packet,
146                    Err(err) => {
147                        warn!(error = %err, "discarding malformed GVSP packet");
148                        continue;
149                    }
150                };
151                match packet {
152                    GvspPacket::Leader { block_id, width, height, pixel_format, timestamp, .. } => {
153                        state = Some(BlockState {
154                            block_id,
155                            width,
156                            height,
157                            pixel_format: PixelFormat::from_code(pixel_format),
158                            timestamp,
159                            payload: BytesMut::new(),
160                        });
161                    }
162                    GvspPacket::Payload { block_id, data, .. } => {
163                        if let Some(active) = state.as_mut()
164                            && active.block_id == block_id {
165                                active.payload.extend_from_slice(data.as_ref());
166                            }
167                    }
168                    GvspPacket::Trailer { block_id, status, chunk_data, .. } => {
169                        let Some(active) = state.take() else { continue };
170                        if active.block_id != block_id {
171                            continue;
172                        }
173                        if status != 0 {
174                            warn!(block_id, status, "trailer reported non-zero status");
175                        }
176                        if !chunk_data.is_empty()
177                            && let Err(err) = parse_chunk_bytes(chunk_data.as_ref()) {
178                                warn!(block_id, error = %err, "failed to decode chunk payload");
179                            }
180                        let frame = Frame {
181                            payload: active.payload.freeze(),
182                            width: active.width,
183                            height: active.height,
184                            pixel_format: active.pixel_format,
185                            chunks: None,
186                            ts_dev: Some(active.timestamp),
187                            ts_host: Some(camera.map_dev_ts(active.timestamp)),
188                        };
189                        let latency = frame
190                            .host_time()
191                            .and_then(|ts| SystemTime::now().duration_since(ts).ok());
192                        stats.record_frame(frame.payload.len(), latency);
193                    }
194                }
195            }
196        }
197    }
198
199    camera.acquisition_stop().context("stop acquisition")?;
200    if interrupted {
201        println!("Benchmark interrupted by user.");
202    }
203
204    let snapshot = stream.stats();
205    println!(
206        "Summary: frames={} bytes={} drops={} resends={} avg_fps={:.1} avg_mbps={:.2}",
207        snapshot.frames,
208        snapshot.bytes,
209        snapshot.drops,
210        snapshot.resends,
211        snapshot.avg_fps,
212        snapshot.avg_mbps,
213    );
214
215    let report = BenchReport {
216        duration_s: duration.as_secs(),
217        frames: snapshot.frames,
218        bytes: snapshot.bytes,
219        avg_fps: snapshot.avg_fps,
220        avg_mbps: snapshot.avg_mbps,
221        drops: snapshot.drops,
222        resends: snapshot.resends,
223        mode: args.mode.clone(),
224    };
225
226    if let Some(path) = args.json_out.as_ref() {
227        let file = File::create(path).with_context(|| format!("create {}", path.display()))?;
228        serde_json::to_writer_pretty(file, &report)
229            .with_context(|| format!("write {}", path.display()))?;
230        info!(file = %path.display(), "wrote benchmark report");
231    }
232
233    if emit_json && args.json_out.is_none() {
234        common::print_json(&report)?;
235    }
236
237    Ok(())
238}
239
240#[derive(Debug, Clone, Copy, PartialEq, Eq)]
241enum StreamMode {
242    Unicast,
243    Multicast,
244}
245
246fn parse_mode(value: &str) -> Result<StreamMode> {
247    match value.to_ascii_lowercase().as_str() {
248        "unicast" => Ok(StreamMode::Unicast),
249        "multicast" => Ok(StreamMode::Multicast),
250        other => bail!("unknown stream mode '{other}' (expected unicast or multicast)"),
251    }
252}