@@ -28,6 +28,8 @@ use arroyo_server_common::wrap_start;
2828use arroyo_types:: { MachineId , PipelineId , WorkerId , from_micros} ;
2929use arroyo_worker:: job_controller:: job_metrics:: JobMetrics ;
3030use cornucopia_async:: DatabaseSource ;
31+ use lazy_static:: lazy_static;
32+ use prometheus:: { IntGaugeVec , register_int_gauge_vec} ;
3133use states:: { Created , State , StateMachine } ;
3234use std:: collections:: { HashMap , HashSet } ;
3335use std:: env;
@@ -51,6 +53,63 @@ pub mod schedulers;
5153mod states;
5254
5355const TTL_PIPELINE_CLEANUP_TIME : Duration = Duration :: from_secs ( 60 * 60 ) ;
56+ const JOB_STATES : [ & str ; 6 ] = [
57+ "running" ,
58+ "transitioning" ,
59+ "stopped" ,
60+ "failed" ,
61+ "user_failed" ,
62+ "unknown" ,
63+ ] ;
64+
65+ lazy_static ! {
66+ static ref JOBS_BY_STATE : IntGaugeVec = register_int_gauge_vec!(
67+ "arroyo_controller_jobs" ,
68+ "Current number of jobs by operational state" ,
69+ & [ "state" ]
70+ )
71+ . unwrap( ) ;
72+ }
73+
74+ fn metric_job_state (
75+ state : Option < & str > ,
76+ running_desired : bool ,
77+ failure_domain : Option < & str > ,
78+ ) -> & ' static str {
79+ match ( state. unwrap_or ( "Created" ) , running_desired) {
80+ ( "Failed" , _) if failure_domain == Some ( "user" ) => "user_failed" ,
81+ ( "Failed" , _) => "failed" ,
82+ ( "Running" , true ) => "running" ,
83+ ( "Created" | "Stopped" | "Finished" , false ) => "stopped" ,
84+ (
85+ "Created" | "Compiling" | "Scheduling" | "Running" | "Rescaling" | "CheckpointStopping"
86+ | "Recovering" | "Restarting" | "Stopping" | "Stopped" | "Finishing" | "Finished"
87+ | "Failing" ,
88+ _,
89+ ) => "transitioning" ,
90+ _ => "unknown" ,
91+ }
92+ }
93+
94+ fn job_state_counts < ' a > (
95+ jobs : impl Iterator < Item = ( Option < & ' a str > , bool , Option < & ' a str > ) > ,
96+ ) -> HashMap < & ' static str , i64 > {
97+ let mut counts = HashMap :: new ( ) ;
98+ for ( state, running_desired, failure_domain) in jobs {
99+ * counts
100+ . entry ( metric_job_state ( state, running_desired, failure_domain) )
101+ . or_default ( ) += 1 ;
102+ }
103+ counts
104+ }
105+
106+ fn update_job_state_metrics ( counts : & HashMap < & ' static str , i64 > ) {
107+ for state in JOB_STATES {
108+ JOBS_BY_STATE
109+ . with_label_values ( & [ state] )
110+ . set ( counts. get ( state) . copied ( ) . unwrap_or_default ( ) ) ;
111+ }
112+ }
54113
55114include ! ( concat!( env!( "OUT_DIR" ) , "/controller-sql.rs" ) ) ;
56115
@@ -645,6 +704,15 @@ impl ControllerServer {
645704 while !token. is_cancelled ( ) {
646705 let client = db. client ( ) . await ?;
647706 let res = queries:: controller_queries:: fetch_all_jobs ( & client) . await ?;
707+ let state_counts = job_state_counts ( res. iter ( ) . map ( |p| {
708+ (
709+ p. state . as_deref ( ) ,
710+ p. stop == StopMode :: none,
711+ p. failure_domain . as_deref ( ) ,
712+ )
713+ } ) ) ;
714+ update_job_state_metrics ( & state_counts) ;
715+
648716 for p in res {
649717 let id = Arc :: new ( p. id ) ;
650718 let config = JobConfig {
@@ -799,3 +867,77 @@ impl ControllerServer {
799867 Ok ( local_addr. port ( ) )
800868 }
801869}
870+
871+ #[ cfg( test) ]
872+ mod tests {
873+ use super :: * ;
874+
875+ #[ test]
876+ fn metric_job_states_use_operational_buckets ( ) {
877+ assert_eq ! ( metric_job_state( Some ( "Running" ) , true , None ) , "running" ) ;
878+ assert_eq ! (
879+ metric_job_state( Some ( "Running" ) , false , None ) ,
880+ "transitioning"
881+ ) ;
882+ assert_eq ! ( metric_job_state( Some ( "Created" ) , false , None ) , "stopped" ) ;
883+ assert_eq ! (
884+ metric_job_state( Some ( "Created" ) , true , None ) ,
885+ "transitioning"
886+ ) ;
887+ assert_eq ! ( metric_job_state( Some ( "Stopped" ) , false , None ) , "stopped" ) ;
888+ assert_eq ! (
889+ metric_job_state( Some ( "Stopped" ) , true , None ) ,
890+ "transitioning"
891+ ) ;
892+ assert_eq ! ( metric_job_state( Some ( "Finished" ) , false , None ) , "stopped" ) ;
893+ assert_eq ! (
894+ metric_job_state( Some ( "Finished" ) , true , None ) ,
895+ "transitioning"
896+ ) ;
897+ assert_eq ! ( metric_job_state( None , false , None ) , "stopped" ) ;
898+
899+ for state in [
900+ "Compiling" ,
901+ "Scheduling" ,
902+ "Rescaling" ,
903+ "CheckpointStopping" ,
904+ "Recovering" ,
905+ "Restarting" ,
906+ "Stopping" ,
907+ "Finishing" ,
908+ "Failing" ,
909+ ] {
910+ assert_eq ! ( metric_job_state( Some ( state) , true , None ) , "transitioning" ) ;
911+ }
912+
913+ assert_eq ! (
914+ metric_job_state( Some ( "Failed" ) , false , Some ( "user" ) ) ,
915+ "user_failed"
916+ ) ;
917+ assert_eq ! (
918+ metric_job_state( Some ( "Failed" ) , true , Some ( "internal" ) ) ,
919+ "failed"
920+ ) ;
921+ assert_eq ! ( metric_job_state( Some ( "Unexpected" ) , true , None ) , "unknown" ) ;
922+ }
923+
924+ #[ test]
925+ fn job_state_metrics_clear_absent_states ( ) {
926+ let counts = job_state_counts (
927+ [
928+ ( Some ( "Running" ) , true , None ) ,
929+ ( Some ( "Running" ) , true , None ) ,
930+ ( Some ( "Unexpected" ) , true , None ) ,
931+ ]
932+ . into_iter ( ) ,
933+ ) ;
934+ update_job_state_metrics ( & counts) ;
935+ assert_eq ! ( JOBS_BY_STATE . with_label_values( & [ "running" ] ) . get( ) , 2 ) ;
936+ assert_eq ! ( JOBS_BY_STATE . with_label_values( & [ "unknown" ] ) . get( ) , 1 ) ;
937+
938+ update_job_state_metrics ( & HashMap :: new ( ) ) ;
939+ for state in JOB_STATES {
940+ assert_eq ! ( JOBS_BY_STATE . with_label_values( & [ state] ) . get( ) , 0 ) ;
941+ }
942+ }
943+ }
0 commit comments