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}