@@ -525,57 +525,6 @@ fn spawn_stdout_reader(stdout: ChildStdout, state: Arc<ClientState>) -> JoinHand
525525 } )
526526}
527527
528- #[ cfg( test) ]
529- pub ( crate ) fn controlled_pending_request (
530- client : & PetJsonRpcClient ,
531- submitted_at : Instant ,
532- ) -> ( PendingRequest , impl FnOnce ( Result < Value , String > , Instant ) ) {
533- let id = REQUEST_ID . fetch_add ( 1 , Ordering :: SeqCst ) ;
534- let ( sender, receiver) = mpsc:: channel ( ) ;
535- client
536- . inner
537- . state
538- . pending
539- . lock ( )
540- . unwrap ( )
541- . insert ( id, sender. clone ( ) ) ;
542- let request = PendingRequest {
543- id,
544- method : "controlled" . to_string ( ) ,
545- submitted_at,
546- receiver,
547- client : client. clone ( ) ,
548- } ;
549- let send = move |result, received_at| {
550- sender
551- . send ( TimedResponse {
552- result,
553- received_at,
554- } )
555- . expect ( "controlled pending receiver should remain connected" ) ;
556- } ;
557- ( request, send)
558- }
559-
560- #[ cfg( test) ]
561- pub ( crate ) fn pending_request_registered (
562- client : & PetJsonRpcClient ,
563- request : & PendingRequest ,
564- ) -> bool {
565- client
566- . inner
567- . state
568- . pending
569- . lock ( )
570- . unwrap ( )
571- . contains_key ( & request. id )
572- }
573-
574- #[ cfg( test) ]
575- pub ( crate ) fn pending_request_count ( client : & PetJsonRpcClient ) -> usize {
576- client. inner . state . pending . lock ( ) . unwrap ( ) . len ( )
577- }
578-
579528fn spawn_stderr_reader (
580529 stderr : impl Read + Send + ' static ,
581530 state : Arc < ClientState > ,
@@ -649,3 +598,113 @@ pub(crate) fn read_message(reader: &mut BufReader<ChildStdout>) -> io::Result<Op
649598 } ) ?;
650599 Ok ( Some ( message) )
651600}
601+
602+ #[ cfg( test) ]
603+ mod tests {
604+ use super :: * ;
605+
606+ fn controlled_pending_request (
607+ client : & PetJsonRpcClient ,
608+ submitted_at : Instant ,
609+ ) -> ( PendingRequest , impl FnOnce ( Result < Value , String > , Instant ) ) {
610+ let id = REQUEST_ID . fetch_add ( 1 , Ordering :: SeqCst ) ;
611+ let ( sender, receiver) = mpsc:: channel ( ) ;
612+ client
613+ . inner
614+ . state
615+ . pending
616+ . lock ( )
617+ . unwrap ( )
618+ . insert ( id, sender. clone ( ) ) ;
619+ let request = PendingRequest {
620+ id,
621+ method : "controlled" . to_string ( ) ,
622+ submitted_at,
623+ receiver,
624+ client : client. clone ( ) ,
625+ } ;
626+ let send = move |result, received_at| {
627+ sender
628+ . send ( TimedResponse {
629+ result,
630+ received_at,
631+ } )
632+ . expect ( "controlled pending receiver should remain connected" ) ;
633+ } ;
634+ ( request, send)
635+ }
636+
637+ fn pending_request_registered ( client : & PetJsonRpcClient , request : & PendingRequest ) -> bool {
638+ client
639+ . inner
640+ . state
641+ . pending
642+ . lock ( )
643+ . unwrap ( )
644+ . contains_key ( & request. id )
645+ }
646+
647+ fn pending_request_count ( client : & PetJsonRpcClient ) -> usize {
648+ client. inner . state . pending . lock ( ) . unwrap ( ) . len ( )
649+ }
650+
651+ #[ test]
652+ fn pending_request_uses_receipt_time_and_operation_deadline ( ) {
653+ let client = PetJsonRpcClient :: spawn ( ) . expect ( "failed to spawn idle PET client" ) ;
654+
655+ let submitted_at = Instant :: now ( ) - Duration :: from_secs ( 10 ) ;
656+ let received_at = submitted_at + Duration :: from_millis ( 125 ) ;
657+ let ( received, send_received) = controlled_pending_request ( & client, submitted_at) ;
658+ send_received ( Ok ( json ! ( { "ok" : true } ) ) , received_at) ;
659+ let ( result, latency) = received
660+ . wait ( Duration :: from_secs ( 30 ) )
661+ . expect ( "recorded response should be returned" ) ;
662+ assert_eq ! ( result, json!( { "ok" : true } ) ) ;
663+ assert_eq ! ( latency, Duration :: from_millis( 125 ) ) ;
664+ assert_eq ! ( pending_request_count( & client) , 0 ) ;
665+
666+ let consumed_after_deadline_at = Instant :: now ( ) - Duration :: from_secs ( 31 ) ;
667+ let ( consumed_after_deadline, send_timely) =
668+ controlled_pending_request ( & client, consumed_after_deadline_at) ;
669+ send_timely (
670+ Ok ( json ! ( { "timely" : true } ) ) ,
671+ consumed_after_deadline_at + Duration :: from_secs ( 29 ) ,
672+ ) ;
673+ let ( result, latency) = consumed_after_deadline
674+ . wait ( Duration :: from_secs ( 30 ) )
675+ . expect ( "a timely received response remains valid when consumed after the deadline" ) ;
676+ assert_eq ! ( result, json!( { "timely" : true } ) ) ;
677+ assert_eq ! ( latency, Duration :: from_secs( 29 ) ) ;
678+
679+ let received_after_deadline_at = Instant :: now ( ) - Duration :: from_secs ( 31 ) ;
680+ let ( received_after_deadline, send_late) =
681+ controlled_pending_request ( & client, received_after_deadline_at) ;
682+ send_late (
683+ Ok ( json ! ( { "late" : true } ) ) ,
684+ received_after_deadline_at + Duration :: from_secs ( 30 ) + Duration :: from_millis ( 1 ) ,
685+ ) ;
686+ let error = received_after_deadline
687+ . wait ( Duration :: from_secs ( 30 ) )
688+ . expect_err ( "a response received after the operation deadline must time out" ) ;
689+ assert ! ( error. contains( "Timed out waiting for controlled response" ) ) ;
690+ assert_eq ! ( pending_request_count( & client) , 0 ) ;
691+
692+ let expired_at = Instant :: now ( ) - Duration :: from_secs ( 31 ) ;
693+ let ( expired, _keep_sender_connected) = controlled_pending_request ( & client, expired_at) ;
694+ assert ! ( pending_request_registered( & client, & expired) ) ;
695+ let error = expired
696+ . wait ( Duration :: from_secs ( 30 ) )
697+ . expect_err ( "an expired operation must not receive a fresh wait budget" ) ;
698+ assert ! ( error. contains( "Timed out waiting for controlled response" ) ) ;
699+ assert_eq ! ( pending_request_count( & client) , 0 ) ;
700+
701+ let ( dropped, _keep_sender_connected) = controlled_pending_request ( & client, Instant :: now ( ) ) ;
702+ assert ! ( pending_request_registered( & client, & dropped) ) ;
703+ drop ( dropped) ;
704+ assert_eq ! (
705+ pending_request_count( & client) ,
706+ 0 ,
707+ "dropping a pending request must remove its map entry"
708+ ) ;
709+ }
710+ }
0 commit comments