Skip to main content

runner_core/
executor.rs

1// Copyright 2026 Kotelnikovekb
2//
3// Licensed under the Apache License, Version 2.0 (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//     http://apache.org
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
15use anyhow::Result;
16use async_trait::async_trait;
17use runner_protocol::{ExecutionRequirements, IsolationLevel, JobResult, JobSpec};
18use serde::{Deserialize, Serialize};
19use std::collections::{BTreeMap, BTreeSet};
20use std::sync::Arc;
21use tokio_util::sync::CancellationToken;
22
23#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
24pub struct ExecutorCapabilities {
25    #[serde(default)]
26    pub isolation: IsolationLevel,
27    #[serde(default)]
28    pub capabilities: BTreeSet<String>,
29}
30
31impl Default for ExecutorCapabilities {
32    fn default() -> Self {
33        Self {
34            isolation: IsolationLevel::Container,
35            capabilities: BTreeSet::new(),
36        }
37    }
38}
39
40impl ExecutorCapabilities {
41    pub fn new(names: impl IntoIterator<Item = impl Into<String>>) -> Self {
42        Self {
43            isolation: IsolationLevel::Container,
44            capabilities: names.into_iter().map(Into::into).collect(),
45        }
46    }
47
48    pub fn supports(&self, name: &str) -> bool {
49        self.capabilities.contains(name)
50    }
51
52    pub fn satisfies(&self, requirements: &ExecutionRequirements) -> Result<()> {
53        if requirements.isolation != IsolationLevel::Container
54            && requirements.isolation != self.isolation
55        {
56            anyhow::bail!(
57                "executor does not support required isolation: {:?}",
58                requirements.isolation
59            )
60        }
61        for capability in &requirements.capabilities {
62            if !self.supports(capability) {
63                anyhow::bail!("executor does not support capability: {capability}")
64            }
65        }
66        Ok(())
67    }
68}
69
70#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
71pub struct DoctorCheck {
72    pub name: String,
73    pub healthy: bool,
74    pub message: String,
75}
76
77#[derive(Debug, Clone, PartialEq, Eq, Serialize, Deserialize)]
78pub struct DoctorReport {
79    pub executor: String,
80    pub healthy: bool,
81    pub capabilities: ExecutorCapabilities,
82    pub checks: Vec<DoctorCheck>,
83}
84
85#[async_trait]
86pub trait Executor: Send + Sync {
87    fn capabilities(&self) -> ExecutorCapabilities;
88
89    async fn doctor(&self) -> Result<DoctorReport>;
90
91    async fn run(&self, spec: JobSpec, cancel: Option<CancellationToken>) -> Result<JobResult>;
92
93    async fn cleanup(&self) -> Result<()>;
94}
95
96#[derive(Default)]
97pub struct ExecutorFactory {
98    executors: BTreeMap<String, Arc<dyn Executor>>,
99}
100
101impl ExecutorFactory {
102    pub fn register(&mut self, name: impl Into<String>, executor: Arc<dyn Executor>) -> Result<()> {
103        let name = name.into();
104        if name.trim().is_empty() {
105            anyhow::bail!("executor name must not be empty")
106        }
107        if self.executors.contains_key(&name) {
108            anyhow::bail!("executor already registered: {name}")
109        }
110        self.executors.insert(name, executor);
111        Ok(())
112    }
113
114    pub fn get(&self, name: &str) -> Result<Arc<dyn Executor>> {
115        self.executors
116            .get(name)
117            .cloned()
118            .ok_or_else(|| anyhow::anyhow!("executor is not registered: {name}"))
119    }
120
121    pub fn select_for(&self, spec: &JobSpec) -> Result<Arc<dyn Executor>> {
122        let executor = self.get(&spec.executor)?;
123        let requirements = spec.execution.clone().unwrap_or(ExecutionRequirements {
124            isolation: IsolationLevel::Container,
125            capabilities: Vec::new(),
126        });
127        executor
128            .capabilities()
129            .satisfies(&requirements)
130            .map_err(|error| anyhow::anyhow!("executor admission failed: {error}"))?;
131        Ok(executor)
132    }
133
134    pub fn names(&self) -> impl Iterator<Item = &str> {
135        self.executors.keys().map(String::as_str)
136    }
137}
138
139#[cfg(test)]
140mod tests {
141    use super::{Executor, ExecutorCapabilities, ExecutorFactory};
142    use crate::mock::MockExecutor;
143    use runner_protocol::{ExecutionRequirements, IsolationLevel};
144    use std::sync::Arc;
145
146    #[test]
147    fn capabilities_are_extensible_and_deterministic() {
148        let capabilities = ExecutorCapabilities::new(["network", "artifacts", "network"]);
149
150        assert!(capabilities.supports("network"));
151        assert_eq!(capabilities.capabilities.len(), 2);
152    }
153
154    #[test]
155    fn capabilities_reject_unsupported_isolation_and_names() {
156        let capabilities = ExecutorCapabilities::new(["artifacts"]);
157
158        assert!(
159            capabilities
160                .satisfies(&ExecutionRequirements {
161                    isolation: IsolationLevel::Container,
162                    capabilities: vec!["artifacts".into()],
163                })
164                .is_ok()
165        );
166        assert!(
167            capabilities
168                .satisfies(&ExecutionRequirements {
169                    isolation: IsolationLevel::Sandboxed,
170                    capabilities: Vec::new(),
171                })
172                .is_err()
173        );
174        assert!(
175            capabilities
176                .satisfies(&ExecutionRequirements {
177                    isolation: IsolationLevel::Container,
178                    capabilities: vec!["unknown".into()],
179                })
180                .is_err()
181        );
182    }
183
184    #[tokio::test]
185    async fn factory_rejects_duplicate_and_unknown_executors() {
186        let executor = Arc::new(MockExecutor);
187        let mut factory = ExecutorFactory::default();
188
189        factory.register("mock", executor.clone()).unwrap();
190        assert!(factory.register("mock", executor).is_err());
191        assert!(factory.get("missing").is_err());
192        assert_eq!(
193            factory.get("mock").unwrap().capabilities(),
194            MockExecutor.capabilities()
195        );
196    }
197}