1+ use bonsai_bt:: Behavior :: Wait ;
2+ use bonsai_bt:: {
3+ Behavior :: Race , Behavior :: Action , Behavior :: Sequence ,
4+ Event , Status , Timer , UpdateArgs , BT , RUNNING ,
5+ } ;
6+ use futures:: FutureExt ;
7+ use std:: collections:: HashMap ;
8+ use std:: sync:: mpsc:: { channel, Receiver } ;
9+ use std:: thread:: sleep;
10+ use std:: time:: Duration ;
11+ use tokio:: time:: sleep as async_sleep;
12+ use rand:: Rng ;
13+
14+
15+ #[ derive( Clone , Debug , serde:: Deserialize , serde:: Serialize ) ]
16+ pub enum MissionAction {
17+ /// The main job that finishes after a random delay
18+ DoWork ,
19+ /// Hard deadline that fires after a fixed delay
20+ OnTimeout ,
21+ }
22+
23+ pub struct MissionState {
24+ pub work : Option < Receiver < Status > > ,
25+ }
26+
27+ /// Simulates a unit of work whose duration is random.
28+ /// Sometimes it finishes before the timeout and sometimes it doesn't.
29+ async fn do_work_task ( tx : std:: sync:: mpsc:: Sender < Status > ) {
30+ let work_ms: u64 = rand:: thread_rng ( ) . gen_range ( 200 ..=1200 ) ;
31+ println ! ( "[do_work] started." ) ;
32+
33+ let step = Duration :: from_millis ( 100 ) ;
34+ let mut elapsed = 0u64 ;
35+ while elapsed < work_ms {
36+
37+ if tx. send ( Status :: Running ) . is_err ( ) {
38+ println ! ( "[do_work] preempted by timeout, stopping." ) ;
39+ return ;
40+ }
41+ async_sleep ( step) . await ;
42+ elapsed += step. as_millis ( ) as u64 ;
43+ }
44+
45+ println ! ( "[do_work] finished after {elapsed} ms" ) ;
46+ let _ = tx. send ( Status :: Success ) ;
47+ }
48+
49+ async fn tick (
50+ timer : & mut Timer ,
51+ state : & mut MissionState ,
52+ bt : & mut BT < MissionAction , HashMap < String , serde_json:: Value > > ,
53+ ) -> std:: option:: Option < ( Status , f64 ) > {
54+ let dt = timer. get_dt ( ) ;
55+ let e: Event = UpdateArgs { dt } . into ( ) ;
56+
57+ bt. tick ( & e, & mut |args : bonsai_bt:: ActionArgs < Event , MissionAction > , _| {
58+ match * args. action {
59+ MissionAction :: DoWork => {
60+ if let Some ( rx) = & state. work {
61+ match rx. recv ( ) {
62+ Ok ( Status :: Running ) => RUNNING ,
63+ Ok ( Status :: Success ) => {
64+ state. work = None ;
65+ ( Status :: Success , args. dt )
66+ }
67+ Ok ( Status :: Failure ) | Err ( _) => {
68+ state. work = None ;
69+ ( Status :: Failure , args. dt )
70+ }
71+ }
72+ } else {
73+ let ( tx, rx) = channel ( ) ;
74+ let ( job, handle) = do_work_task ( tx) . remote_handle ( ) ;
75+ handle. forget ( ) ;
76+ tokio:: spawn ( job) ;
77+ state. work = Some ( rx) ;
78+ match state. work . as_ref ( ) . unwrap ( ) . recv ( ) . unwrap ( ) {
79+ Status :: Running => RUNNING ,
80+ s => ( s, args. dt ) ,
81+ }
82+ }
83+ }
84+
85+ MissionAction :: OnTimeout => {
86+ eprintln ! ( "do_work timed out!" ) ;
87+ ( Status :: Failure , args. dt )
88+ }
89+ }
90+ } )
91+ }
92+
93+ #[ tokio:: main]
94+ async fn main ( ) {
95+ const TIMEOUT_S : f64 = 0.6 ;
96+
97+ let behavior = Sequence ( vec ! [
98+ Race ( vec![
99+ Action ( MissionAction :: DoWork ) ,
100+ Sequence ( vec![
101+ Wait ( TIMEOUT_S ) ,
102+ Action ( MissionAction :: OnTimeout )
103+ ] )
104+ ] ) ,
105+ ] ) ;
106+
107+ let mut bt = BT :: new ( behavior, HashMap :: new ( ) ) ;
108+ let mut timer = Timer :: init_time ( ) ;
109+ let mut state = MissionState {
110+ work : None ,
111+ } ;
112+
113+ loop {
114+ sleep ( Duration :: from_millis ( 50 ) ) ;
115+ if tick ( & mut timer, & mut state, & mut bt) . await . is_none ( ) {
116+ break ;
117+ }
118+ }
119+ }
0 commit comments