@@ -20,15 +20,12 @@ use docopt::Docopt;
2020use std:: env;
2121use std:: fs;
2222use std:: fs:: File ;
23- use std:: io:: { self , BufRead , BufReader , Write , Read } ;
23+ use std:: io:: { self , BufRead , BufReader , Write } ;
2424use std:: path:: Path ;
2525use std:: sync:: Arc ;
2626use std:: sync:: Mutex ;
2727use std:: thread;
2828use std:: time:: { Duration , SystemTime , UNIX_EPOCH } ;
29- use std:: net:: { TcpListener , TcpStream , SocketAddr } ;
30- use std:: str:: FromStr ;
31- use std:: io:: ErrorKind ;
3229
3330// This is a simple app that pairs with the Secluso camera, receives motion videos,
3431// and launches livestream sessions.
@@ -41,7 +38,6 @@ use std::io::ErrorKind;
4138const CAMERA_ADDR : & str = "127.0.0.1" ;
4239const CAMERA_NAME : & str = "Camera" ;
4340const DATA_DIR : & str = "example_app_data" ;
44- const FIRST_APP_ADDR : & str = "127.0.0.1" ;
4541
4642pub const MAX_ALLOWED_MSG_LEN : u64 = 65536 ;
4743
@@ -94,6 +90,11 @@ fn main() -> io::Result<()> {
9490 let clients: Arc < Mutex < Option < Box < Clients > > > > = Arc :: new ( Mutex :: new ( None ) ) ;
9591 let http_client = HttpClient :: new ( server_addr, server_username, server_password) ;
9692
93+ // We assume here that the new secret is shared via
94+ // another channel, e.g., QR code scan.
95+ // Also, a new secret needs to be used for every app added.
96+ let add_app_secret = vec ! [ 2u8 ; NUM_SECRET_BYTES ] ;
97+
9798 if first_time {
9899 if args. flag_reset {
99100 panic ! ( "No state to reset!" ) ;
@@ -122,25 +123,19 @@ fn main() -> io::Result<()> {
122123 )
123124 } else {
124125 println ! ( "Sending the add_app request" ) ;
125- let addr = SocketAddr :: from_str ( & ( FIRST_APP_ADDR . to_owned ( ) + ":12350" ) )
126- . map_err ( |e| io:: Error :: other ( format ! ( "{e}" ) ) ) ?;
127-
128- let mut stream = TcpStream :: connect ( & addr) ?;
129126
130127 // get key packages
131128 let key_packages_vec = get_key_packages ( & mut clients. lock ( ) . unwrap ( ) ) ?;
132129
133- write_varying_len ( & mut stream, & key_packages_vec) ?;
134-
135- let new_app_data_vec = read_varying_len ( & mut stream) ?;
136-
137- // We assume here that the new secret is shared via
138- // another channel, e.g., QR code scan.
139- let new_secret = vec ! [ 2u8 ; NUM_SECRET_BYTES ] ;
130+ println ! ( "About to send add_app request" ) ;
131+ http_client. add_app_request ( "test_add_app_request_token" , key_packages_vec) ?;
132+ println ! ( "About to wait for add_app response" ) ;
133+ let new_app_data_vec = http_client. add_app_check ( "test_add_app_response_token" ) ?;
134+ println ! ( "Received add_app response" ) ;
140135
141136 let epochs: [ u64 ; NUM_MLS_CLIENTS ] = join_camera_groups (
142137 & mut clients. lock ( ) . unwrap ( ) ,
143- new_secret . clone ( ) ,
138+ add_app_secret . clone ( ) ,
144139 new_app_data_vec,
145140 ) ?;
146141
@@ -166,30 +161,35 @@ fn main() -> io::Result<()> {
166161 }
167162 }
168163
169- let add_app_request: Arc < Mutex < Option < TcpStream > > > = Arc :: new ( Mutex :: new ( None ) ) ;
164+ let add_app_request: Arc < Mutex < Option < Vec < u8 > > > > = Arc :: new ( Mutex :: new ( None ) ) ;
170165
171166 if !args. flag_secondary_app {
172167 let add_app_request_clone = Arc :: clone ( & add_app_request) ;
168+ let http_client_clone = http_client. clone ( ) ;
173169
174170 thread:: spawn ( move || loop {
175- let listener = TcpListener :: bind ( "0.0.0.0:12350" ) . unwrap ( ) ;
176- for incoming in listener. incoming ( ) {
177- match incoming {
178- Ok ( stream) => {
179- println ! ( "Incoming connection accepted." ) ;
180- let mut stream_opt = add_app_request_clone. lock ( ) . unwrap ( ) ;
181- * stream_opt = Some ( stream) ;
182- } ,
183-
184- Err ( e) => {
185- println ! ( "Incoming connection error: {e}" ) ;
186- }
187- }
188- }
171+ println ! ( "About to wait for add_app request" ) ;
172+ match http_client_clone. add_app_check ( "test_add_app_request_token" ) {
173+ Ok ( data) => {
174+ println ! ( "Received add_app request." ) ;
175+ let mut data_opt = add_app_request_clone. lock ( ) . unwrap ( ) ;
176+ * data_opt = Some ( data) ;
177+ } ,
178+
179+ Err ( e) => {
180+ println ! ( "Error listening for add_app requests: {e}" ) ;
181+ }
182+ }
189183 } ) ;
190184 }
191185
192- main_loop ( clients, http_client, add_app_request, args. flag_num_iters ) ?;
186+ main_loop (
187+ clients,
188+ http_client,
189+ add_app_request,
190+ args. flag_num_iters ,
191+ add_app_secret,
192+ ) ?;
193193
194194 Ok ( ( ) )
195195}
@@ -212,8 +212,9 @@ fn deregister_all(
212212fn main_loop (
213213 clients : Arc < Mutex < Option < Box < Clients > > > > ,
214214 http_client : HttpClient ,
215- add_app_request : Arc < Mutex < Option < TcpStream > > > ,
215+ add_app_request : Arc < Mutex < Option < Vec < u8 > > > > ,
216216 num_iters : usize ,
217+ add_app_secret : Vec < u8 > ,
217218) -> io:: Result < ( ) > {
218219 for iter in 0 ..num_iters {
219220 thread:: sleep ( Duration :: from_secs ( 1 ) ) ;
@@ -230,15 +231,16 @@ fn main_loop(
230231 livestream ( Arc :: clone ( & clients) , & http_client, 2 ) ?;
231232 }
232233
233- let mut add_app_stream_opt = add_app_request. lock ( ) . unwrap ( ) ;
234- if let Some ( add_app_stream ) = add_app_stream_opt . as_mut ( ) {
234+ let mut add_app_data_opt = add_app_request. lock ( ) . unwrap ( ) ;
235+ if let Some ( add_app_data ) = add_app_data_opt . as_ref ( ) {
235236 println ! ( "Add app request detected" ) ;
236237 handle_add_app_request (
237238 Arc :: clone ( & clients) ,
238239 & http_client,
239- add_app_stream,
240+ add_app_data,
241+ add_app_secret. clone ( ) ,
240242 ) ?;
241- * add_app_stream_opt = None ;
243+ * add_app_data_opt = None ;
242244 }
243245 }
244246
@@ -248,15 +250,15 @@ fn main_loop(
248250fn handle_add_app_request (
249251 clients : Arc < Mutex < Option < Box < Clients > > > > ,
250252 http_client : & HttpClient ,
251- stream : & mut TcpStream ,
253+ add_app_data : & Vec < u8 > ,
254+ add_app_secret : Vec < u8 > ,
252255) -> io:: Result < ( ) > {
253256 println ! ( "handle_add_app_request called" ) ;
254- let new_secret = vec ! [ 2u8 ; NUM_SECRET_BYTES ] ;
255257
256- let new_app_key_packages_vec = read_varying_len ( stream ) ? ;
258+ let new_app_key_packages_vec = add_app_data . clone ( ) ;
257259
258260 let config_msg_enc =
259- generate_add_app_request_config_command ( & mut clients. lock ( ) . unwrap ( ) , new_app_key_packages_vec, new_secret . clone ( ) ) ?;
261+ generate_add_app_request_config_command ( & mut clients. lock ( ) . unwrap ( ) , new_app_key_packages_vec, add_app_secret . clone ( ) ) ?;
260262
261263 let config_group_name = get_group_name ( & mut clients. lock ( ) . unwrap ( ) , "config" ) ?;
262264
@@ -286,13 +288,13 @@ fn handle_add_app_request(
286288 let new_app_data_vec = process_add_app_config_response (
287289 & mut clients. lock ( ) . unwrap ( ) ,
288290 config_response. clone ( ) ,
289- new_secret ,
291+ add_app_secret ,
290292 ) . unwrap ( ) ;
291293
292294 increment_epoch ( "motion_epoch" ) ;
293295 increment_epoch ( "thumbnail_epoch" ) ;
294296
295- write_varying_len ( stream , & new_app_data_vec) ?;
297+ http_client . add_app_request ( "test_add_app_response_token" , new_app_data_vec) ?;
296298
297299 Ok ( ( ) )
298300}
@@ -516,66 +518,3 @@ fn fetch_livestream_chunk(
516518 format ! ( "Error: could not fetch livestream chunk (timeout)!" ) ,
517519 ) ) ;
518520}
519-
520- // FIXME: copied from camera_hub/src/pairing.rs.
521- fn write_varying_len ( stream : & mut TcpStream , msg : & [ u8 ] ) -> io:: Result < ( ) > {
522- // FIXME: is u64 necessary?
523- let len = msg. len ( ) as u64 ;
524- let len_data = len. to_be_bytes ( ) ;
525-
526- stream. write_all ( & len_data) ?;
527- stream. write_all ( msg) ?;
528- stream. flush ( ) ?;
529-
530- Ok ( ( ) )
531- }
532-
533- fn read_varying_len ( stream : & mut TcpStream ) -> io:: Result < Vec < u8 > > {
534- let mut len_data = [ 0u8 ; 8 ] ;
535-
536- match stream. read_exact ( & mut len_data) {
537- Ok ( _) => { }
538- Err ( ref e) if e. kind ( ) == ErrorKind :: WouldBlock => {
539- return Err ( io:: Error :: new (
540- ErrorKind :: WouldBlock ,
541- "Length read would block" ,
542- ) ) ;
543- }
544- Err ( e) => return Err ( e) ,
545- }
546-
547- let len = u64:: from_be_bytes ( len_data) ;
548-
549- if len > MAX_ALLOWED_MSG_LEN {
550- println ! ( "Communicated message length ({len}) exceeds the allowed length ({MAX_ALLOWED_MSG_LEN})" ) ;
551- return Err ( io:: Error :: new (
552- ErrorKind :: InvalidInput ,
553- "Intended message length is too large" ,
554- ) )
555- }
556-
557- let mut msg = vec ! [ 0u8 ; len as usize ] ;
558- let mut offset = 0 ;
559-
560- while offset < msg. len ( ) {
561- match stream. read ( & mut msg[ offset..] ) {
562- Ok ( 0 ) => {
563- return Err ( io:: Error :: new (
564- ErrorKind :: UnexpectedEof ,
565- "Socket closed during read" ,
566- ) )
567- }
568- Ok ( n) => {
569- offset += n;
570- }
571- Err ( ref e) if e. kind ( ) == ErrorKind :: WouldBlock => {
572- // retry a few times with a short delay
573- thread:: sleep ( Duration :: from_millis ( 10 ) ) ;
574- continue ;
575- }
576- Err ( e) => return Err ( e) ,
577- }
578- }
579-
580- Ok ( msg)
581- }
0 commit comments