1use 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}