wowlab_sentinel/scheduler/
assignment.rs1use uuid::Uuid;
4use wowlab_types::sim::FastSet;
5
6use super::{
7 backlog::NodeBacklogs,
8 dispatch::{
9 AssignmentClock, AssignmentPublisher, ClaimTokenSource, DispatchDependencies, Dispatcher,
10 StateAssignmentPublisher, SystemClock, UuidClaimTokenSource,
11 },
12 eligibility::EligibilityContext,
13 planning::{JobPlan, build_runtime, order_for_dispatch, plan_jobs},
14 repository::{AssignmentRepository, RepositoryError, SqlAssignmentRepository},
15};
16use crate::{state::RuntimeState, utils::filter_refresh::FilterMap};
17
18#[derive(Debug, thiserror::Error)]
19pub(super) enum AssignmentError {
20 #[error(transparent)]
21 Repository(#[from] RepositoryError),
22}
23
24#[derive(Debug, Default, Eq, PartialEq)]
25pub(super) struct AssignmentRound {
26 pub jobs: usize,
27 pub planned: usize,
28 pub started: usize,
29 pub assigned: usize,
30 pub invalid_configs: usize,
31 pub no_eligible_node: usize,
32 pub no_work: usize,
33 pub rejected_publications: usize,
34 pub no_online_nodes: bool,
35}
36
37pub(super) struct AssignmentWorkflow<'a, R, P, C, T> {
38 runtime: &'a RuntimeState,
39 filters: &'a FilterMap,
40 repository: R,
41 publisher: P,
42 clock: C,
43 claim_tokens: T,
44}
45
46#[cfg(test)]
47pub(super) struct AssignmentDependencies<R, P, C, T> {
48 pub repository: R,
49 pub publisher: P,
50 pub clock: C,
51 pub claim_tokens: T,
52}
53
54pub(super) fn assignment_workflow(
55 state: &crate::state::ServerState,
56) -> AssignmentWorkflow<
57 '_,
58 SqlAssignmentRepository<'_>,
59 StateAssignmentPublisher<'_>,
60 SystemClock,
61 UuidClaimTokenSource,
62> {
63 AssignmentWorkflow {
64 runtime: &state.runtime,
65 filters: &state.filters,
66 repository: SqlAssignmentRepository::new(
67 state.dbs.get::<super::SchedulerDb>(),
68 state.config.scheduler_batch_size,
69 ),
70 publisher: StateAssignmentPublisher::new(state),
71 clock: SystemClock,
72 claim_tokens: UuidClaimTokenSource,
73 }
74}
75
76impl<R, P, C, T> AssignmentWorkflow<'_, R, P, C, T>
77where
78 R: AssignmentRepository,
79 P: AssignmentPublisher,
80 C: AssignmentClock,
81 T: ClaimTokenSource,
82{
83 #[cfg(test)]
84 pub(super) fn new<'a>(
85 runtime: &'a RuntimeState,
86 filters: &'a FilterMap,
87 dependencies: AssignmentDependencies<R, P, C, T>,
88 ) -> AssignmentWorkflow<'a, R, P, C, T> {
89 AssignmentWorkflow {
90 runtime,
91 filters,
92 repository: dependencies.repository,
93 publisher: dependencies.publisher,
94 clock: dependencies.clock,
95 claim_tokens: dependencies.claim_tokens,
96 }
97 }
98
99 pub(super) async fn run(&self) -> Result<AssignmentRound, AssignmentError> {
100 let jobs = self.repository.pending_jobs().await?;
101
102 if jobs.is_empty() {
103 tracing::debug!("No pending jobs found");
104 metrics::gauge!(crate::telemetry::CHUNKS_PENDING).set(0.0);
105
106 return Ok(AssignmentRound::default());
107 }
108
109 let pending: Vec<_> = jobs
110 .iter()
111 .map(|job| (job.id, job.status.as_str()))
112 .collect();
113
114 tracing::info!(?pending, "Pending jobs");
115 let job_count = jobs.len();
116 let mut plan = plan_jobs(jobs);
117 let started = self.initialize_runtimes(&plan).await?;
118
119 log_plan(&plan, started);
120
121 let nodes = self.repository.online_nodes().await?;
122
123 if nodes.is_empty() {
124 tracing::info!("No online nodes available for assignment");
125
126 return Ok(AssignmentRound {
127 jobs: job_count,
128 planned: plan.jobs.len(),
129 started,
130 invalid_configs: plan.invalid_configs,
131 no_online_nodes: true,
132 ..AssignmentRound::default()
133 });
134 }
135
136 let public_keys = nodes
137 .iter()
138 .map(|node| node.public_key.clone())
139 .collect::<Vec<_>>();
140 let permissions = self.repository.permissions(&public_keys).await?;
141 let mut backlogs = {
142 let runtimes = self.runtime.jobs.read().await;
143
144 NodeBacklogs::from_runtimes(&runtimes)
145 };
146
147 tracing::debug!(?backlogs, "Node backlogs");
148 let user_ids = unique_user_ids(&plan);
149 let user_discord_ids = self.repository.user_discord_ids(&user_ids).await?;
150 let friend_memberships = self.repository.friend_memberships(&user_ids).await?;
151
152 order_for_dispatch(&mut plan.jobs);
153 let eligibility = EligibilityContext {
154 permissions: &permissions,
155 filters: self.filters,
156 friend_memberships: &friend_memberships,
157 user_discord_ids: &user_discord_ids,
158 };
159 let dispatch = Dispatcher::new(
160 self.runtime,
161 DispatchDependencies {
162 publisher: &self.publisher,
163 clock: &self.clock,
164 claim_tokens: &self.claim_tokens,
165 },
166 )
167 .dispatch(&plan.jobs, &nodes, &eligibility, &mut backlogs)
168 .await;
169
170 log_dispatch(&dispatch, plan.invalid_configs);
171 record_assignment_metrics(dispatch.assigned);
172 tracing::info!(
173 assigned = dispatch.assigned,
174 jobs = job_count,
175 "Assignment round complete"
176 );
177
178 Ok(AssignmentRound {
179 jobs: job_count,
180 planned: plan.jobs.len(),
181 started,
182 assigned: dispatch.assigned,
183 invalid_configs: plan.invalid_configs,
184 no_eligible_node: dispatch.no_eligible_node,
185 no_work: dispatch.no_work,
186 rejected_publications: dispatch.rejected_publications,
187 no_online_nodes: false,
188 })
189 }
190
191 async fn initialize_runtimes(&self, plan: &JobPlan) -> Result<usize, AssignmentError> {
192 let mut started = 0usize;
193
194 for planned in &plan.jobs {
195 let already_tracked = self.runtime.jobs.read().await.contains(&planned.job.id);
196
197 if already_tracked {
198 continue;
199 }
200
201 self.runtime
202 .jobs
203 .write()
204 .await
205 .insert(build_runtime(&planned.job, &planned.config));
206 let rows = self.repository.mark_running(planned.job.id).await?;
207
208 started += usize::try_from(rows).unwrap_or(usize::MAX);
209 }
210
211 Ok(started)
212 }
213}
214
215impl AssignmentError {
216 pub(super) const fn is_pending_fetch(&self) -> bool {
217 matches!(self, Self::Repository(RepositoryError::FetchPending(_)))
218 }
219}
220
221#[async_trait::async_trait]
222impl<T> AssignmentPublisher for &T
223where
224 T: AssignmentPublisher + ?Sized,
225{
226 async fn job_running(&self, job_id: Uuid) -> super::dispatch::PublishOutcome {
227 (**self).job_running(job_id).await
228 }
229
230 async fn chunk(
231 &self,
232 public_key: &wowlab_common::NodePublicKey,
233 payload: &wowlab_common::RuntimeChunkPayload,
234 ) -> super::dispatch::PublishOutcome {
235 (**self).chunk(public_key, payload).await
236 }
237}
238
239impl<T> AssignmentClock for &T
240where
241 T: AssignmentClock + ?Sized,
242{
243 fn now_ms(&self) -> u64 {
244 (**self).now_ms()
245 }
246}
247
248impl<T> ClaimTokenSource for &T
249where
250 T: ClaimTokenSource + ?Sized,
251{
252 fn next(&self) -> wowlab_common::ClaimToken {
253 (**self).next()
254 }
255}
256
257fn unique_user_ids(plan: &JobPlan) -> Vec<Uuid> {
258 let unique: FastSet<Uuid> = plan
259 .jobs
260 .iter()
261 .map(|planned| planned.job.user_id)
262 .collect();
263
264 unique.into_iter().collect()
265}
266
267fn log_plan(plan: &JobPlan, started: usize) {
268 if plan.invalid_configs > 0 {
269 tracing::error!(
270 invalid_configs = plan.invalid_configs,
271 "Rejected jobs with invalid sentinel config"
272 );
273 }
274
275 tracing::info!(
276 planned = plan.jobs.len(),
277 started_jobs = started,
278 "Planned assignment round"
279 );
280}
281
282fn log_dispatch(report: &super::dispatch::DispatchReport, invalid_configs: usize) {
283 tracing::info!(
284 assigned_count = report.assigned,
285 no_eligible_node = report.no_eligible_node,
286 no_work = report.no_work,
287 invalid_configs,
288 rejected_publications = report.rejected_publications,
289 "Assignment round distribution"
290 );
291}
292
293fn record_assignment_metrics(assigned: usize) {
294 if assigned == 0 {
295 return;
296 }
297
298 metrics::counter!(crate::telemetry::CHUNKS_ASSIGNED).increment(assigned as u64);
299 #[expect(
300 clippy::cast_precision_loss,
301 reason = "assignment counts remain exactly representable by operational metrics"
302 )]
303 metrics::gauge!(crate::telemetry::CHUNKS_RUNNING).increment(assigned as f64);
304}
305
306#[cfg(test)]
307mod tests;