Skip to main content

wowlab_sentinel/scheduler/
assignment.rs

1//! Application workflow for one scheduler assignment round.
2
3use 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;