@@ -8,14 +8,63 @@ use cap_media_info::VideoInfo;
88use futures:: { FutureExt , channel:: mpsc, future:: BoxFuture } ;
99use std:: sync:: {
1010 Arc ,
11- atomic:: { AtomicBool , Ordering } ,
11+ atomic:: { AtomicBool , AtomicU64 , Ordering } ,
1212} ;
1313use tokio:: sync:: oneshot;
1414
15+ struct CameraFrameScaler {
16+ context : ffmpeg:: software:: scaling:: Context ,
17+ source_width : u32 ,
18+ source_height : u32 ,
19+ target_width : u32 ,
20+ target_height : u32 ,
21+ }
22+
23+ impl CameraFrameScaler {
24+ fn new (
25+ src_width : u32 ,
26+ src_height : u32 ,
27+ src_format : ffmpeg:: format:: Pixel ,
28+ dst_width : u32 ,
29+ dst_height : u32 ,
30+ dst_format : ffmpeg:: format:: Pixel ,
31+ ) -> anyhow:: Result < Self > {
32+ let context = ffmpeg:: software:: scaling:: Context :: get (
33+ src_format,
34+ src_width,
35+ src_height,
36+ dst_format,
37+ dst_width,
38+ dst_height,
39+ ffmpeg:: software:: scaling:: Flags :: BILINEAR ,
40+ ) ?;
41+
42+ Ok ( Self {
43+ context,
44+ source_width : src_width,
45+ source_height : src_height,
46+ target_width : dst_width,
47+ target_height : dst_height,
48+ } )
49+ }
50+
51+ fn matches_source ( & self , width : u32 , height : u32 ) -> bool {
52+ self . source_width == width && self . source_height == height
53+ }
54+
55+ fn scale ( & mut self , input : & ffmpeg:: frame:: Video ) -> anyhow:: Result < ffmpeg:: frame:: Video > {
56+ let mut output = ffmpeg:: frame:: Video :: empty ( ) ;
57+ self . context . run ( input, & mut output) ?;
58+ output. set_pts ( input. pts ( ) ) ;
59+ Ok ( output)
60+ }
61+ }
62+
1563pub struct Camera {
1664 feed_lock : Arc < CameraFeedLock > ,
1765 stop_tx : Option < oneshot:: Sender < ( ) > > ,
1866 stopped : Arc < AtomicBool > ,
67+ original_video_info : VideoInfo ,
1968}
2069
2170impl VideoSource for Camera {
@@ -32,6 +81,11 @@ impl VideoSource for Camera {
3281 {
3382 let ( tx, rx) = flume:: bounded ( 256 ) ;
3483
84+ let original_video_info = * feed_lock. video_info ( ) ;
85+ let original_width = original_video_info. width ;
86+ let original_height = original_video_info. height ;
87+ let original_format = original_video_info. pixel_format ;
88+
3589 feed_lock
3690 . ask ( camera:: AddSender ( tx) )
3791 . await
@@ -40,15 +94,22 @@ impl VideoSource for Camera {
4094 let ( stop_tx, stop_rx) = oneshot:: channel ( ) ;
4195 let stopped = Arc :: new ( AtomicBool :: new ( false ) ) ;
4296 let stopped_clone = stopped. clone ( ) ;
97+ let scaled_frame_count = Arc :: new ( AtomicU64 :: new ( 0 ) ) ;
98+ let scaled_count_clone = scaled_frame_count. clone ( ) ;
4399
44100 tokio:: spawn ( async move {
45- tracing:: debug!( "Camera source task started" ) ;
101+ tracing:: debug!(
102+ original_width,
103+ original_height,
104+ "Camera source task started"
105+ ) ;
46106 let mut frame_count: u64 = 0 ;
47107 let mut sent_count: u64 = 0 ;
48108 let mut dropped_count: u64 = 0 ;
49109 let start = std:: time:: Instant :: now ( ) ;
50110 let mut video_tx = video_tx;
51111 let mut stop_rx = stop_rx. fuse ( ) ;
112+ let mut scaler: Option < CameraFrameScaler > = None ;
52113
53114 loop {
54115 if stopped_clone. load ( Ordering :: Relaxed ) {
@@ -64,17 +125,88 @@ impl VideoSource for Camera {
64125 }
65126 result = rx. recv_async( ) => {
66127 match result {
67- Ok ( frame) => {
128+ Ok ( mut frame) => {
68129 frame_count += 1 ;
130+
131+ let frame_width = frame. inner. width( ) ;
132+ let frame_height = frame. inner. height( ) ;
133+
134+ if frame_width != original_width || frame_height != original_height {
135+ let needs_new_scaler = scaler
136+ . as_ref( )
137+ . map_or( true , |s| !s. matches_source( frame_width, frame_height) ) ;
138+
139+ if needs_new_scaler {
140+ let frame_format = frame. inner. format( ) ;
141+ match CameraFrameScaler :: new(
142+ frame_width,
143+ frame_height,
144+ frame_format,
145+ original_width,
146+ original_height,
147+ original_format,
148+ ) {
149+ Ok ( new_scaler) => {
150+ tracing:: info!(
151+ src_width = frame_width,
152+ src_height = frame_height,
153+ dst_width = original_width,
154+ dst_height = original_height,
155+ "Camera source: created scaler for dimension change"
156+ ) ;
157+ scaler = Some ( new_scaler) ;
158+ }
159+ Err ( e) => {
160+ tracing:: warn!(
161+ src_width = frame_width,
162+ src_height = frame_height,
163+ error = %e,
164+ "Camera source: failed to create scaler, dropping frame"
165+ ) ;
166+ dropped_count += 1 ;
167+ continue ;
168+ }
169+ }
170+ }
171+
172+ if let Some ( s) = & mut scaler {
173+ match s. scale( & frame. inner) {
174+ Ok ( scaled) => {
175+ frame. inner = scaled;
176+ scaled_count_clone. fetch_add( 1 , Ordering :: Relaxed ) ;
177+ }
178+ Err ( e) => {
179+ if dropped_count. is_multiple_of( 30 ) {
180+ tracing:: warn!(
181+ error = %e,
182+ "Camera source: scale failed, dropping frame"
183+ ) ;
184+ }
185+ dropped_count += 1 ;
186+ continue ;
187+ }
188+ }
189+ }
190+ } else if scaler. is_some( ) {
191+ let total_scaled = scaled_count_clone. load( Ordering :: Relaxed ) ;
192+ tracing:: info!(
193+ total_scaled,
194+ "Camera source: dimensions restored to original, removing scaler"
195+ ) ;
196+ scaler = None ;
197+ }
198+
69199 match video_tx. try_send( frame) {
70200 Ok ( ( ) ) => {
71201 sent_count += 1 ;
72- if sent_count. is_multiple_of( 30 ) {
202+ if sent_count. is_multiple_of( 300 ) {
203+ let total_scaled = scaled_count_clone. load( Ordering :: Relaxed ) ;
73204 tracing:: debug!(
74- "Camera source: sent {} frames, dropped {} in {:?}" ,
75205 sent_count,
76206 dropped_count,
77- start. elapsed( )
207+ total_scaled,
208+ elapsed = ?start. elapsed( ) ,
209+ "Camera source stats"
78210 ) ;
79211 }
80212 }
@@ -83,15 +215,15 @@ impl VideoSource for Camera {
83215 dropped_count += 1 ;
84216 if dropped_count. is_multiple_of( 30 ) {
85217 tracing:: warn!(
86- "Camera source: encoder can't keep up, dropped {} frames so far" ,
87- dropped_count
218+ dropped_count ,
219+ "Camera source: encoder can't keep up"
88220 ) ;
89221 }
90222 } else if e. is_disconnected( ) {
91223 tracing:: debug!(
92- "Camera source: pipeline closed after {} sent, {} dropped" ,
93224 sent_count,
94- dropped_count
225+ dropped_count,
226+ "Camera source: pipeline closed"
95227 ) ;
96228 break ;
97229 }
@@ -100,9 +232,10 @@ impl VideoSource for Camera {
100232 }
101233 Err ( e) => {
102234 tracing:: debug!(
103- "Camera feed disconnected (rx closed) after {} frames in {:?}: {e}" ,
104235 frame_count,
105- start. elapsed( )
236+ elapsed = ?start. elapsed( ) ,
237+ error = %e,
238+ "Camera feed disconnected (rx closed)"
106239 ) ;
107240 break ;
108241 }
@@ -113,24 +246,27 @@ impl VideoSource for Camera {
113246
114247 drop ( video_tx) ;
115248
249+ let total_scaled = scaled_count_clone. load ( Ordering :: Relaxed ) ;
116250 tracing:: info!(
117- "Camera source finished: {} received, {} sent, {} dropped in {:?}" ,
118251 frame_count,
119252 sent_count,
120253 dropped_count,
121- start. elapsed( )
254+ total_scaled,
255+ elapsed = ?start. elapsed( ) ,
256+ "Camera source finished"
122257 ) ;
123258 } ) ;
124259
125260 Ok ( Self {
126261 feed_lock,
127262 stop_tx : Some ( stop_tx) ,
128263 stopped,
264+ original_video_info,
129265 } )
130266 }
131267
132268 fn video_info ( & self ) -> VideoInfo {
133- * self . feed_lock . video_info ( )
269+ self . original_video_info
134270 }
135271
136272 fn stop ( & mut self ) -> BoxFuture < ' _ , anyhow:: Result < ( ) > > {
0 commit comments