@@ -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;
@@ -52,6 +54,43 @@ mod states;
5254
5355const TTL_PIPELINE_CLEANUP_TIME : Duration = Duration :: from_secs ( 60 * 60 ) ;
5456
57+ lazy_static ! {
58+ static ref JOBS_BY_STATE : IntGaugeVec = register_int_gauge_vec!(
59+ "arroyo_controller_jobs" ,
60+ "Current number of jobs by controller state" ,
61+ & [ "state" ]
62+ )
63+ . unwrap( ) ;
64+ }
65+
66+ fn metric_job_state < ' a > ( state : Option < & ' a str > , failure_domain : Option < & str > ) -> & ' a str {
67+ let state = state. unwrap_or ( "Created" ) ;
68+ if state == "Failed" && failure_domain == Some ( "user" ) {
69+ "UserFailed"
70+ } else {
71+ state
72+ }
73+ }
74+
75+ fn job_state_counts < ' a > (
76+ jobs : impl Iterator < Item = ( Option < & ' a str > , Option < & ' a str > ) > ,
77+ ) -> HashMap < & ' a str , i64 > {
78+ let mut counts = HashMap :: new ( ) ;
79+ for ( state, failure_domain) in jobs {
80+ * counts
81+ . entry ( metric_job_state ( state, failure_domain) )
82+ . or_default ( ) += 1 ;
83+ }
84+ counts
85+ }
86+
87+ fn update_job_state_metrics ( counts : & HashMap < & str , i64 > ) {
88+ JOBS_BY_STATE . reset ( ) ;
89+ for ( state, count) in counts {
90+ JOBS_BY_STATE . with_label_values ( & [ state] ) . set ( * count) ;
91+ }
92+ }
93+
5594include ! ( concat!( env!( "OUT_DIR" ) , "/controller-sql.rs" ) ) ;
5695
5796use crate :: schedulers:: { ManualScheduler , NodeScheduler , ProcessScheduler , Scheduler } ;
@@ -645,6 +684,12 @@ impl ControllerServer {
645684 while !token. is_cancelled ( ) {
646685 let client = db. client ( ) . await ?;
647686 let res = queries:: controller_queries:: fetch_all_jobs ( & client) . await ?;
687+ let state_counts = job_state_counts (
688+ res. iter ( )
689+ . map ( |p| ( p. state . as_deref ( ) , p. failure_domain . as_deref ( ) ) ) ,
690+ ) ;
691+ update_job_state_metrics ( & state_counts) ;
692+
648693 for p in res {
649694 let id = Arc :: new ( p. id ) ;
650695 let config = JobConfig {
@@ -799,3 +844,39 @@ impl ControllerServer {
799844 Ok ( local_addr. port ( ) )
800845 }
801846}
847+
848+ #[ cfg( test) ]
849+ mod tests {
850+ use prometheus:: core:: Collector ;
851+
852+ use super :: * ;
853+
854+ #[ test]
855+ fn metric_job_states_preserve_raw_states ( ) {
856+ for state in [ "Created" , "Running" , "Finished" , "Unexpected" ] {
857+ assert_eq ! ( metric_job_state( Some ( state) , None ) , state) ;
858+ }
859+ assert_eq ! ( metric_job_state( None , None ) , "Created" ) ;
860+
861+ assert_eq ! ( metric_job_state( Some ( "Failed" ) , Some ( "user" ) ) , "UserFailed" ) ;
862+ assert_eq ! ( metric_job_state( Some ( "Failed" ) , Some ( "internal" ) ) , "Failed" ) ;
863+ }
864+
865+ #[ test]
866+ fn job_state_metrics_clear_absent_states ( ) {
867+ let counts = job_state_counts (
868+ [
869+ ( Some ( "Running" ) , None ) ,
870+ ( Some ( "Running" ) , None ) ,
871+ ( Some ( "Failed" ) , Some ( "user" ) ) ,
872+ ]
873+ . into_iter ( ) ,
874+ ) ;
875+ update_job_state_metrics ( & counts) ;
876+ assert_eq ! ( JOBS_BY_STATE . with_label_values( & [ "Running" ] ) . get( ) , 2 ) ;
877+ assert_eq ! ( JOBS_BY_STATE . with_label_values( & [ "UserFailed" ] ) . get( ) , 1 ) ;
878+
879+ update_job_state_metrics ( & HashMap :: new ( ) ) ;
880+ assert ! ( JOBS_BY_STATE . collect( ) [ 0 ] . get_metric( ) . is_empty( ) ) ;
881+ }
882+ }
0 commit comments