Skip to main content

solti_model/domain/
capability.rs

1//! # Agent capabilities
2//!
3//! [`AgentCapabilities`] is an immutable runner capability snapshot.
4//! [`RunnerCapability`] describes one registered runner.
5//!
6//! Runner order is preserved.
7//! Workload GVKs inside each runner are stored in canonical order.
8
9use std::collections::HashSet;
10
11use serde::{Deserialize, Deserializer, Serialize};
12
13use crate::{Labels, ModelError, ModelResult, WORKLOAD_API_VERSION, WorkloadTypeMeta, validation};
14
15const EMBEDDED_WORKLOAD_KIND: &str = "Embedded";
16
17/// One registered runner and the workload GVKs it can execute.
18///
19/// Labels are the same static labels used by `runnerSelector`.
20#[derive(Clone, Debug, Eq, PartialEq, Serialize)]
21#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
22#[cfg_attr(feature = "schema", schemars(deny_unknown_fields))]
23#[serde(rename_all = "camelCase")]
24pub struct RunnerCapability {
25    #[cfg_attr(
26        feature = "schema",
27        schemars(schema_with = "crate::schema::runner_name")
28    )]
29    name: String,
30    labels: Labels,
31    #[cfg_attr(
32        feature = "schema",
33        schemars(schema_with = "crate::schema::runner_workload_types")
34    )]
35    workload_types: Vec<WorkloadTypeMeta>,
36}
37
38impl<'de> Deserialize<'de> for RunnerCapability {
39    fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
40    where
41        D: Deserializer<'de>,
42    {
43        #[derive(Deserialize)]
44        #[serde(rename_all = "camelCase", deny_unknown_fields)]
45        struct RawRunnerCapability {
46            name: String,
47            #[serde(default)]
48            labels: Labels,
49            workload_types: Vec<WorkloadTypeMeta>,
50        }
51
52        let raw = RawRunnerCapability::deserialize(deserializer)?;
53        Self::new(raw.name, raw.labels, raw.workload_types).map_err(serde::de::Error::custom)
54    }
55}
56
57impl RunnerCapability {
58    /// Creates a registered runner capability.
59    ///
60    /// Workload types are stored in canonical GVK order.
61    ///
62    /// # Errors
63    ///
64    /// Returns [`ModelError::Invalid`] when the name or labels are invalid, no workload is declared, a GVK is duplicated, or Embedded is declared.
65    pub fn new(
66        name: impl Into<String>,
67        labels: Labels,
68        mut workload_types: Vec<WorkloadTypeMeta>,
69    ) -> ModelResult<Self> {
70        let name = name.into();
71        if name.is_empty() {
72            return Err(ModelError::Invalid(
73                "runner capability name must not be empty".into(),
74            ));
75        }
76        validation::validate_label_value("runner capability name", &name)?;
77        labels.validate()?;
78        if workload_types.is_empty() {
79            return Err(ModelError::Invalid(
80                "runner capability must declare at least one workload GVK".into(),
81            ));
82        }
83        if workload_types.iter().any(|workload| {
84            workload.api_version() == WORKLOAD_API_VERSION
85                && workload.kind() == EMBEDDED_WORKLOAD_KIND
86        }) {
87            return Err(ModelError::Invalid(
88                "runner capability must not declare the Embedded workload".into(),
89            ));
90        }
91
92        workload_types.sort_by(|left, right| {
93            left.api_version()
94                .cmp(right.api_version())
95                .then_with(|| left.kind().cmp(right.kind()))
96        });
97        if workload_types.windows(2).any(|pair| pair[0] == pair[1]) {
98            return Err(ModelError::Invalid(
99                "runner capability contains a duplicate workload GVK".into(),
100            ));
101        }
102
103        Ok(Self {
104            name,
105            labels,
106            workload_types,
107        })
108    }
109
110    /// Registered runner name.
111    #[inline]
112    pub fn name(&self) -> &str {
113        &self.name
114    }
115
116    /// Static labels used by `runnerSelector`.
117    #[inline]
118    pub fn labels(&self) -> &Labels {
119        &self.labels
120    }
121
122    /// Canonically ordered workload GVKs handled by the runner.
123    #[inline]
124    pub fn workload_types(&self) -> &[WorkloadTypeMeta] {
125        &self.workload_types
126    }
127}
128
129/// Immutable snapshot of agent execution capabilities.
130#[derive(Clone, Debug, Default, Eq, PartialEq, Serialize)]
131#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
132#[cfg_attr(feature = "schema", schemars(deny_unknown_fields))]
133#[serde(rename_all = "camelCase")]
134pub struct AgentCapabilities {
135    runners: Vec<RunnerCapability>,
136}
137
138impl<'de> Deserialize<'de> for AgentCapabilities {
139    fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
140    where
141        D: Deserializer<'de>,
142    {
143        #[derive(Deserialize)]
144        #[serde(rename_all = "camelCase", deny_unknown_fields)]
145        struct RawAgentCapabilities {
146            #[serde(default)]
147            runners: Vec<RunnerCapability>,
148        }
149
150        let raw = RawAgentCapabilities::deserialize(deserializer)?;
151        Self::new(raw.runners).map_err(serde::de::Error::custom)
152    }
153}
154
155impl AgentCapabilities {
156    /// Creates a capability snapshot in runner registration order.
157    ///
158    /// # Errors
159    ///
160    /// Returns [`ModelError::Invalid`] when runner names are duplicated.
161    pub fn new(runners: Vec<RunnerCapability>) -> ModelResult<Self> {
162        let mut names = HashSet::with_capacity(runners.len());
163        if runners
164            .iter()
165            .any(|runner| !names.insert(runner.name().to_owned()))
166        {
167            return Err(ModelError::Invalid(
168                "agent capabilities contain a duplicate runner name".into(),
169            ));
170        }
171        Ok(Self { runners })
172    }
173
174    /// Registered runners in routing priority order.
175    #[inline]
176    pub fn runners(&self) -> &[RunnerCapability] {
177        &self.runners
178    }
179}
180
181#[cfg(test)]
182mod tests {
183    use super::*;
184
185    fn workload(api_version: &str, kind: &str) -> WorkloadTypeMeta {
186        WorkloadTypeMeta::new(api_version, kind).unwrap()
187    }
188
189    #[test]
190    fn capability_canonicalizes_workload_types() {
191        let capability = RunnerCapability::new(
192            "runner-1",
193            Labels::new(),
194            vec![
195                workload("tasks.example.io/v1", "Resize"),
196                workload(WORKLOAD_API_VERSION, "Subprocess"),
197            ],
198        )
199        .unwrap();
200
201        assert_eq!(
202            capability
203                .workload_types()
204                .iter()
205                .map(|workload| (workload.api_version(), workload.kind()))
206                .collect::<Vec<_>>(),
207            vec![
208                ("solti.io/v1", "Subprocess"),
209                ("tasks.example.io/v1", "Resize"),
210            ]
211        );
212    }
213
214    #[test]
215    fn capability_rejects_invalid_registration_data() {
216        assert!(
217            RunnerCapability::new(
218                "",
219                Labels::new(),
220                vec![workload("solti.io/v1", "Subprocess")]
221            )
222            .is_err()
223        );
224        assert!(
225            RunnerCapability::new(
226                "invalid/name",
227                Labels::new(),
228                vec![workload("solti.io/v1", "Subprocess")],
229            )
230            .is_err()
231        );
232        assert!(RunnerCapability::new("runner", Labels::new(), Vec::new()).is_err());
233        assert!(
234            RunnerCapability::new(
235                "runner",
236                Labels::new(),
237                vec![
238                    workload("solti.io/v1", "Subprocess"),
239                    workload("solti.io/v1", "Subprocess"),
240                ],
241            )
242            .is_err()
243        );
244        assert!(
245            RunnerCapability::new(
246                "runner",
247                Labels::new(),
248                vec![workload(WORKLOAD_API_VERSION, EMBEDDED_WORKLOAD_KIND)],
249            )
250            .is_err()
251        );
252    }
253
254    #[test]
255    fn runner_name_uses_kubernetes_label_value_rules() {
256        let capability = |name: String| {
257            RunnerCapability::new(
258                name,
259                Labels::new(),
260                vec![workload(WORKLOAD_API_VERSION, "Subprocess")],
261            )
262        };
263
264        assert!(capability("Runner_A.1".into()).is_ok());
265        for invalid in [
266            String::new(),
267            "-runner".into(),
268            "runner-".into(),
269            "runner/name".into(),
270            "r".repeat(64),
271        ] {
272            assert!(capability(invalid).is_err());
273        }
274    }
275
276    #[test]
277    fn capabilities_reject_duplicate_runner_names() {
278        let first = RunnerCapability::new(
279            "runner",
280            Labels::new(),
281            vec![workload(WORKLOAD_API_VERSION, "Subprocess")],
282        )
283        .unwrap();
284        let second = RunnerCapability::new(
285            "runner",
286            Labels::new(),
287            vec![workload("tasks.example.io/v1", "Resize")],
288        )
289        .unwrap();
290
291        assert!(AgentCapabilities::new(vec![first, second]).is_err());
292    }
293
294    #[test]
295    fn serde_is_strict_and_validated() {
296        let valid = serde_json::json!({
297            "runners": [{
298                "name": "runner",
299                "labels": {"zone": "eu"},
300                "workloadTypes": [{
301                    "apiVersion": "solti.io/v1",
302                    "kind": "Subprocess"
303                }]
304            }]
305        });
306        let capabilities: AgentCapabilities = serde_json::from_value(valid).unwrap();
307        assert_eq!(capabilities.runners()[0].labels().get("zone"), Some("eu"));
308
309        let unknown = serde_json::json!({
310            "runners": [],
311            "unknown": true
312        });
313        assert!(serde_json::from_value::<AgentCapabilities>(unknown).is_err());
314    }
315}