Coverage Report

Created: 2026-07-21 15:28

next uncovered line (L), next uncovered region (R), next uncovered branch (B)
/build/source/nativelink-scheduler/src/default_scheduler_factory.rs
Line
Count
Source
1
// Copyright 2024 The NativeLink Authors. All rights reserved.
2
//
3
// Licensed under the Functional Source License, Version 1.1, Apache 2.0 Future License (the "License");
4
// you may not use this file except in compliance with the License.
5
// You may obtain a copy of the License at
6
//
7
//    See LICENSE file for details
8
//
9
// Unless required by applicable law or agreed to in writing, software
10
// distributed under the License is distributed on an "AS IS" BASIS,
11
// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12
// See the License for the specific language governing permissions and
13
// limitations under the License.
14
15
use std::sync::Arc;
16
use std::time::SystemTime;
17
18
use nativelink_config::schedulers::{
19
    ExperimentalSimpleSchedulerBackend, SchedulerSpec, SimpleSpec,
20
};
21
use nativelink_config::stores::EvictionPolicy;
22
use nativelink_error::{Error, ResultExt, make_input_err};
23
use nativelink_proto::com::github::trace_machina::nativelink::events::OriginEvent;
24
use nativelink_store::redis_store::{RedisStore, StandardRedisManager};
25
use nativelink_store::store_manager::StoreManager;
26
use nativelink_util::instant_wrapper::InstantWrapper;
27
use redis::aio::ConnectionManager;
28
use tokio::sync::{Notify, mpsc};
29
30
use crate::cache_lookup_scheduler::CacheLookupScheduler;
31
use crate::grpc_scheduler::GrpcScheduler;
32
use crate::historical_resource_scheduler::HistoricalResourceScheduler;
33
use crate::known_platform_property_provider::KnownPlatformPropertyProvider;
34
use crate::memory_awaited_action_db::MemoryAwaitedActionDb;
35
use crate::property_modifier_scheduler::PropertyModifierScheduler;
36
use crate::simple_scheduler::SimpleScheduler;
37
use crate::store_awaited_action_db::StoreAwaitedActionDb;
38
use crate::worker_scheduler::WorkerScheduler;
39
40
/// Default timeout for recently completed actions in seconds.
41
/// If this changes, remember to change the documentation in the config.
42
const DEFAULT_RETAIN_COMPLETED_FOR_S: u32 = 60;
43
44
pub type SchedulerFactoryResults = (
45
    Option<Arc<dyn KnownPlatformPropertyProvider>>,
46
    Option<Arc<dyn WorkerScheduler>>,
47
);
48
49
0
pub async fn scheduler_factory(
50
0
    spec: &SchedulerSpec,
51
0
    store_manager: &StoreManager,
52
0
    maybe_origin_event_tx: Option<&mpsc::Sender<OriginEvent>>,
53
0
) -> Result<SchedulerFactoryResults, Error> {
54
    inner_scheduler_factory(spec, store_manager, maybe_origin_event_tx).await
55
}
56
57
0
async fn inner_scheduler_factory(
58
0
    spec: &SchedulerSpec,
59
0
    store_manager: &StoreManager,
60
0
    maybe_origin_event_tx: Option<&mpsc::Sender<OriginEvent>>,
61
0
) -> Result<SchedulerFactoryResults, Error> {
62
0
    let scheduler: SchedulerFactoryResults = match spec {
63
0
        SchedulerSpec::Simple(spec) => {
64
0
            simple_scheduler_factory(spec, store_manager, SystemTime::now, maybe_origin_event_tx)
65
0
                .await?
66
        }
67
0
        SchedulerSpec::Grpc(spec) => (Some(Arc::new(GrpcScheduler::new(spec)?)), None),
68
0
        SchedulerSpec::CacheLookup(spec) => {
69
0
            let ac_store = store_manager
70
0
                .get_store(&spec.ac_store)
71
0
                .err_tip(|| format!("'ac_store': '{}' does not exist", spec.ac_store))?;
72
0
            let (action_scheduler, worker_scheduler) = Box::pin(inner_scheduler_factory(
73
0
                &spec.scheduler,
74
0
                store_manager,
75
0
                maybe_origin_event_tx,
76
0
            ))
77
0
            .await
78
0
            .err_tip(|| "In nested CacheLookupScheduler construction")?;
79
0
            let cache_lookup_scheduler = Arc::new(CacheLookupScheduler::new(
80
0
                ac_store,
81
0
                action_scheduler.err_tip(|| "Nested scheduler is not an action scheduler")?,
82
0
            )?);
83
0
            (Some(cache_lookup_scheduler), worker_scheduler)
84
        }
85
0
        SchedulerSpec::PropertyModifier(spec) => {
86
0
            let (action_scheduler, worker_scheduler) = Box::pin(inner_scheduler_factory(
87
0
                &spec.scheduler,
88
0
                store_manager,
89
0
                maybe_origin_event_tx,
90
0
            ))
91
0
            .await
92
0
            .err_tip(|| "In nested PropertyModifierScheduler construction")?;
93
0
            let property_modifier_scheduler = Arc::new(PropertyModifierScheduler::new(
94
0
                spec,
95
0
                action_scheduler.err_tip(|| "Nested scheduler is not an action scheduler")?,
96
            ));
97
0
            (Some(property_modifier_scheduler), worker_scheduler)
98
        }
99
0
        SchedulerSpec::HistoricalResource(spec) => {
100
0
            let (action_scheduler, worker_scheduler) = Box::pin(inner_scheduler_factory(
101
0
                &spec.scheduler,
102
0
                store_manager,
103
0
                maybe_origin_event_tx,
104
0
            ))
105
0
            .await
106
0
            .err_tip(|| "In nested HistoricalResourceScheduler construction")?;
107
0
            let historical_resource_scheduler = Arc::new(HistoricalResourceScheduler::new(
108
0
                spec,
109
0
                action_scheduler.err_tip(|| "Nested scheduler is not an action scheduler")?,
110
            ));
111
0
            (Some(historical_resource_scheduler), worker_scheduler)
112
        }
113
    };
114
115
0
    Ok(scheduler)
116
0
}
117
118
0
async fn simple_scheduler_factory(
119
0
    spec: &SimpleSpec,
120
0
    store_manager: &StoreManager,
121
0
    now_fn: fn() -> SystemTime,
122
0
    maybe_origin_event_tx: Option<&mpsc::Sender<OriginEvent>>,
123
0
) -> Result<SchedulerFactoryResults, Error> {
124
0
    match spec
125
0
        .experimental_backend
126
0
        .as_ref()
127
0
        .unwrap_or(&ExperimentalSimpleSchedulerBackend::Memory)
128
    {
129
        ExperimentalSimpleSchedulerBackend::Memory => {
130
0
            let task_change_notify = Arc::new(Notify::new());
131
0
            let awaited_action_db = memory_awaited_action_db_factory(
132
0
                spec.retain_completed_for_s,
133
0
                &task_change_notify,
134
                SystemTime::now,
135
            );
136
0
            let (action_scheduler, worker_scheduler) = SimpleScheduler::new(
137
0
                spec,
138
0
                awaited_action_db,
139
0
                task_change_notify,
140
0
                maybe_origin_event_tx.cloned(),
141
0
            );
142
0
            Ok((Some(action_scheduler), Some(worker_scheduler)))
143
        }
144
0
        ExperimentalSimpleSchedulerBackend::Redis(redis_config) => {
145
0
            let store = store_manager
146
0
                .get_store(redis_config.redis_store.as_ref())
147
0
                .err_tip(|| {
148
0
                    format!(
149
                        "'redis_store': '{}' does not exist",
150
                        redis_config.redis_store
151
                    )
152
0
                })?;
153
0
            let task_change_notify = Arc::new(Notify::new());
154
0
            let store = store
155
0
                .into_inner()
156
0
                .as_any_arc()
157
0
                .downcast::<RedisStore<ConnectionManager, StandardRedisManager<ConnectionManager>>>(
158
                )
159
0
                .map_err(|_| {
160
0
                    make_input_err!(
161
                        "Could not downcast to redis store in RedisAwaitedActionDb::new"
162
                    )
163
0
                })?;
164
0
            let awaited_action_db = StoreAwaitedActionDb::new(
165
0
                store,
166
0
                task_change_notify.clone(),
167
0
                now_fn,
168
0
                Default::default,
169
0
                spec.retain_completed_for_s,
170
0
            )
171
0
            .await
172
0
            .err_tip(|| "In state_manager_factory::redis_state_manager")?;
173
0
            let (action_scheduler, worker_scheduler) = SimpleScheduler::new(
174
0
                spec,
175
0
                awaited_action_db,
176
0
                task_change_notify,
177
0
                maybe_origin_event_tx.cloned(),
178
0
            );
179
0
            Ok((Some(action_scheduler), Some(worker_scheduler)))
180
        }
181
    }
182
0
}
183
184
28
pub fn memory_awaited_action_db_factory<I, NowFn>(
185
28
    mut retain_completed_for_s: u32,
186
28
    task_change_notify: &Arc<Notify>,
187
28
    now_fn: NowFn,
188
28
) -> MemoryAwaitedActionDb<I, NowFn>
189
28
where
190
28
    I: InstantWrapper,
191
28
    NowFn: Fn() -> I + Clone + Send + Sync + 'static,
192
{
193
28
    if retain_completed_for_s == 0 {
194
28
        retain_completed_for_s = DEFAULT_RETAIN_COMPLETED_FOR_S;
195
28
    
}0
196
28
    MemoryAwaitedActionDb::new(
197
28
        &EvictionPolicy {
198
28
            max_seconds: retain_completed_for_s,
199
28
            ..Default::default()
200
28
        },
201
28
        task_change_notify.clone(),
202
28
        now_fn,
203
    )
204
28
}