harn_vm/triggers/worker_queue/
state.rs1use std::collections::BTreeSet;
2
3use serde::{Deserialize, Serialize};
4
5use super::{WorkerQueueClaimHandle, WorkerQueueJob, WorkerQueueResponseRecord};
6use crate::triggers::scheduler::{self, SchedulableJob, SchedulerPolicy, SchedulerState};
7
8#[derive(Clone, Debug, PartialEq, Eq, Serialize, Deserialize)]
9pub struct WorkerQueueSummary {
10 pub queue: String,
11 pub ready: usize,
12 pub in_flight: usize,
13 pub acked: usize,
14 pub purged: usize,
15 pub responses: usize,
16 pub oldest_unclaimed_age_ms: Option<u64>,
17}
18
19#[derive(Clone, Debug, PartialEq, Serialize, Deserialize)]
20pub struct WorkerQueueJobState {
21 pub job_event_id: u64,
22 pub enqueued_at_ms: i64,
23 pub job: WorkerQueueJob,
24 pub active_claim: Option<WorkerQueueClaimHandle>,
25 pub acked: bool,
26 pub purged: bool,
27}
28
29impl WorkerQueueJobState {
30 pub fn is_ready(&self) -> bool {
31 !self.acked && !self.purged && self.active_claim.is_none()
32 }
33}
34
35#[derive(Clone, Debug, Default, PartialEq, Serialize, Deserialize)]
36pub struct WorkerQueueState {
37 pub queue: String,
38 pub responses: Vec<WorkerQueueResponseRecord>,
39 pub jobs: Vec<WorkerQueueJobState>,
40}
41
42impl WorkerQueueState {
43 pub fn summary(&self, now_ms: i64) -> WorkerQueueSummary {
44 let ready = self.jobs.iter().filter(|job| job.is_ready()).count();
45 let in_flight = self
46 .jobs
47 .iter()
48 .filter(|job| !job.acked && !job.purged && job.active_claim.is_some())
49 .count();
50 let acked = self.jobs.iter().filter(|job| job.acked).count();
51 let purged = self.jobs.iter().filter(|job| job.purged).count();
52 let oldest_unclaimed_age_ms = self
53 .jobs
54 .iter()
55 .filter(|job| job.is_ready())
56 .map(|job| now_ms.saturating_sub(job.enqueued_at_ms).max(0) as u64)
57 .max();
58 WorkerQueueSummary {
59 queue: self.queue.clone(),
60 ready,
61 in_flight,
62 acked,
63 purged,
64 responses: self.responses.len(),
65 oldest_unclaimed_age_ms,
66 }
67 }
68
69 pub(super) fn next_ready_job_with_scheduler(
72 &self,
73 scheduler_state: &mut SchedulerState,
74 policy: &SchedulerPolicy,
75 now_ms: i64,
76 excluded_job_event_ids: &BTreeSet<u64>,
77 ) -> Option<(&WorkerQueueJobState, scheduler::SchedulerSelection)> {
78 let candidates: Vec<&WorkerQueueJobState> = self
79 .jobs
80 .iter()
81 .filter(|job| job.is_ready() && !excluded_job_event_ids.contains(&job.job_event_id))
82 .collect();
83 if candidates.is_empty() {
84 return None;
85 }
86 let views: Vec<SchedulableJob<'_>> = candidates
87 .iter()
88 .map(|state| SchedulableJob::from_state(state))
89 .collect();
90
91 let in_flight = scheduler::in_flight_by_key(&self.jobs, policy);
92 scheduler_state.replace_in_flight(in_flight);
93
94 let pick = scheduler_state.select(&views, policy, now_ms)?;
95 candidates
96 .into_iter()
97 .find(|job| job.job_event_id == pick.job_event_id)
98 .map(|job| (job, pick))
99 }
100
101 pub(super) fn active_claim_for(&self, job_event_id: u64) -> Option<&WorkerQueueClaimHandle> {
102 self.jobs
103 .iter()
104 .find(|job| job.job_event_id == job_event_id)
105 .and_then(|job| job.active_claim.as_ref())
106 }
107}