use std::path::PathBuf;
use serde::{Deserialize, Deserializer, Serialize, Serializer};
use serde_json::Value;
use crate::{Flag, ModelError, ModelResult, SubprocessMode, TaskEnv, validation};
pub const WORKLOAD_API_VERSION: &str = "solti.io/v1";
#[derive(Clone, Debug, Eq, Hash, PartialEq, Serialize)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[cfg_attr(feature = "schema", schemars(deny_unknown_fields))]
#[serde(rename_all = "camelCase")]
pub struct WorkloadTypeMeta {
#[cfg_attr(
feature = "schema",
schemars(schema_with = "crate::schema::crd_api_version")
)]
api_version: String,
#[cfg_attr(feature = "schema", schemars(schema_with = "crate::schema::crd_kind"))]
kind: String,
}
impl<'de> Deserialize<'de> for WorkloadTypeMeta {
fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
where
D: Deserializer<'de>,
{
#[derive(Deserialize)]
#[serde(rename_all = "camelCase", deny_unknown_fields)]
struct RawWorkloadTypeMeta {
api_version: String,
kind: String,
}
let raw = RawWorkloadTypeMeta::deserialize(deserializer)?;
Self::new(raw.api_version, raw.kind).map_err(serde::de::Error::custom)
}
}
impl WorkloadTypeMeta {
pub fn new(api_version: impl Into<String>, kind: impl Into<String>) -> ModelResult<Self> {
let type_meta = Self {
api_version: api_version.into(),
kind: kind.into(),
};
type_meta.validate()?;
Ok(type_meta)
}
#[inline]
pub fn api_version(&self) -> &str {
&self.api_version
}
#[inline]
pub fn kind(&self) -> &str {
&self.kind
}
fn validate(&self) -> ModelResult<()> {
validation::validate_crd_api_version("workload apiVersion", &self.api_version)?;
validation::validate_crd_kind("workload kind", &self.kind)
}
}
#[derive(Clone, Debug, Eq, PartialEq, Serialize)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[cfg_attr(feature = "schema", schemars(deny_unknown_fields))]
#[serde(rename_all = "camelCase")]
pub struct EmbeddedSpec {
#[cfg_attr(
feature = "schema",
schemars(schema_with = "crate::schema::non_empty_string")
)]
revision: String,
}
impl<'de> Deserialize<'de> for EmbeddedSpec {
fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
where
D: Deserializer<'de>,
{
#[derive(Deserialize)]
#[serde(deny_unknown_fields)]
struct RawEmbeddedSpec {
revision: String,
}
let raw = RawEmbeddedSpec::deserialize(deserializer)?;
Self::new(raw.revision).map_err(serde::de::Error::custom)
}
}
impl EmbeddedSpec {
pub fn new(revision: impl Into<String>) -> ModelResult<Self> {
let spec = Self {
revision: revision.into(),
};
spec.validate()?;
Ok(spec)
}
#[inline]
pub fn revision(&self) -> &str {
&self.revision
}
fn validate(&self) -> ModelResult<()> {
if self.revision.trim().is_empty() {
return Err(ModelError::Invalid(
"embedded workload revision must not be empty".into(),
));
}
Ok(())
}
}
#[derive(Clone, Debug, Eq, PartialEq)]
#[non_exhaustive]
pub enum TaskWorkload {
Subprocess(SubprocessSpec),
Wasm(WasmSpec),
Container(ContainerSpec),
Embedded(EmbeddedSpec),
Extension(ExtensionWorkload),
}
#[cfg(feature = "schema")]
impl schemars::JsonSchema for TaskWorkload {
fn schema_name() -> std::borrow::Cow<'static, str> {
"TaskWorkload".into()
}
fn json_schema(generator: &mut schemars::SchemaGenerator) -> schemars::Schema {
let subprocess =
workload_envelope_schema("Subprocess", generator.subschema_for::<SubprocessSpec>());
let wasm = workload_envelope_schema("Wasm", generator.subschema_for::<WasmSpec>());
let container =
workload_envelope_schema("Container", generator.subschema_for::<ContainerSpec>());
let embedded =
workload_envelope_schema("Embedded", generator.subschema_for::<EmbeddedSpec>());
let extension = generator.subschema_for::<ExtensionWorkload>();
schemars::json_schema!({
"description": "Kubernetes-style workload GVK and desired state.",
"oneOf": [subprocess, wasm, container, embedded, extension]
})
}
}
#[cfg(feature = "schema")]
fn workload_envelope_schema(kind: &'static str, spec: schemars::Schema) -> schemars::Schema {
schemars::json_schema!({
"type": "object",
"additionalProperties": false,
"required": ["apiVersion", "kind", "spec"],
"properties": {
"apiVersion": {
"type": "string",
"const": WORKLOAD_API_VERSION
},
"kind": {
"type": "string",
"const": kind
},
"spec": spec
}
})
}
#[derive(Clone, Debug, Eq, PartialEq, Serialize)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[cfg_attr(feature = "schema", schemars(deny_unknown_fields))]
#[serde(rename_all = "camelCase")]
pub struct ExtensionWorkload {
#[cfg_attr(
feature = "schema",
schemars(schema_with = "crate::schema::extension_api_version")
)]
api_version: String,
#[cfg_attr(feature = "schema", schemars(schema_with = "crate::schema::crd_kind"))]
kind: String,
#[cfg_attr(
feature = "schema",
schemars(schema_with = "crate::schema::json_object")
)]
spec: Value,
}
impl<'de> Deserialize<'de> for ExtensionWorkload {
fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
where
D: Deserializer<'de>,
{
let raw = RawWorkloadEnvelope::deserialize(deserializer)?;
Self::new(raw.api_version, raw.kind, raw.spec).map_err(serde::de::Error::custom)
}
}
impl ExtensionWorkload {
pub fn new(
api_version: impl Into<String>,
kind: impl Into<String>,
spec: Value,
) -> ModelResult<Self> {
let workload = Self {
api_version: api_version.into(),
kind: kind.into(),
spec,
};
workload.validate()?;
Ok(workload)
}
#[inline]
pub fn api_version(&self) -> &str {
&self.api_version
}
#[inline]
pub fn kind(&self) -> &str {
&self.kind
}
#[inline]
pub fn spec(&self) -> &Value {
&self.spec
}
fn validate(&self) -> ModelResult<()> {
let group = validation::validate_crd_api_version(
"extension workload apiVersion",
&self.api_version,
)?;
validation::validate_crd_kind("extension workload kind", &self.kind)?;
if group == "solti.io" {
return Err(ModelError::Invalid(
format!(
"extension workload GVK {}/{} uses the reserved solti.io API group",
self.api_version, self.kind
)
.into(),
));
}
if !self.spec.is_object() {
return Err(ModelError::Invalid(
"extension workload spec must be a JSON object".into(),
));
}
Ok(())
}
}
impl TaskWorkload {
#[inline]
pub fn api_version(&self) -> &str {
match self {
Self::Extension(workload) => workload.api_version(),
_ => WORKLOAD_API_VERSION,
}
}
pub fn type_meta(&self) -> WorkloadTypeMeta {
WorkloadTypeMeta {
api_version: self.api_version().to_owned(),
kind: self.kind().to_owned(),
}
}
#[inline]
pub fn kind(&self) -> &str {
match self {
Self::Subprocess(_) => "Subprocess",
Self::Container(_) => "Container",
Self::Embedded(_) => "Embedded",
Self::Wasm(_) => "Wasm",
Self::Extension(workload) => workload.kind(),
}
}
pub fn validate(&self) -> ModelResult<()> {
match self {
Self::Subprocess(spec) => spec.mode.validate(),
Self::Container(spec) => spec.validate(),
Self::Wasm(spec) => spec.validate(),
Self::Embedded(spec) => spec.validate(),
Self::Extension(workload) => workload.validate(),
}
}
}
#[derive(Serialize)]
#[serde(rename_all = "camelCase")]
struct WorkloadEnvelope<'a, T> {
api_version: &'a str,
kind: &'a str,
spec: T,
}
impl Serialize for TaskWorkload {
fn serialize<S>(&self, serializer: S) -> Result<S::Ok, S::Error>
where
S: Serializer,
{
match self {
Self::Subprocess(spec) => WorkloadEnvelope {
api_version: self.api_version(),
kind: self.kind(),
spec,
}
.serialize(serializer),
Self::Wasm(spec) => WorkloadEnvelope {
api_version: self.api_version(),
kind: self.kind(),
spec,
}
.serialize(serializer),
Self::Container(spec) => WorkloadEnvelope {
api_version: self.api_version(),
kind: self.kind(),
spec,
}
.serialize(serializer),
Self::Embedded(spec) => WorkloadEnvelope {
api_version: self.api_version(),
kind: self.kind(),
spec,
}
.serialize(serializer),
Self::Extension(workload) => WorkloadEnvelope {
api_version: workload.api_version(),
kind: workload.kind(),
spec: workload.spec(),
}
.serialize(serializer),
}
}
}
#[derive(Deserialize)]
#[serde(rename_all = "camelCase", deny_unknown_fields)]
struct RawWorkloadEnvelope {
api_version: String,
kind: String,
spec: Value,
}
impl<'de> Deserialize<'de> for TaskWorkload {
fn deserialize<D>(deserializer: D) -> Result<Self, D::Error>
where
D: Deserializer<'de>,
{
let raw = RawWorkloadEnvelope::deserialize(deserializer)?;
let workload = if raw.api_version == WORKLOAD_API_VERSION {
match raw.kind.as_str() {
"Subprocess" => Self::Subprocess(
serde_json::from_value(raw.spec).map_err(serde::de::Error::custom)?,
),
"Wasm" => {
Self::Wasm(serde_json::from_value(raw.spec).map_err(serde::de::Error::custom)?)
}
"Container" => Self::Container(
serde_json::from_value(raw.spec).map_err(serde::de::Error::custom)?,
),
"Embedded" => Self::Embedded(
serde_json::from_value(raw.spec).map_err(serde::de::Error::custom)?,
),
_ => Self::Extension(
ExtensionWorkload::new(raw.api_version, raw.kind, raw.spec)
.map_err(serde::de::Error::custom)?,
),
}
} else {
Self::Extension(
ExtensionWorkload::new(raw.api_version, raw.kind, raw.spec)
.map_err(serde::de::Error::custom)?,
)
};
workload.validate().map_err(serde::de::Error::custom)?;
Ok(workload)
}
}
impl WasmSpec {
pub fn new(module: PathBuf, args: Vec<String>, env: TaskEnv) -> Self {
Self { module, args, env }
}
pub fn validate(&self) -> ModelResult<()> {
if self.module.as_os_str().is_empty() {
return Err(ModelError::Invalid(
"wasm module path cannot be empty".into(),
));
}
Ok(())
}
}
impl ContainerSpec {
pub fn new(
image: String,
command: Option<Vec<String>>,
args: Vec<String>,
env: TaskEnv,
) -> Self {
Self {
image,
command,
args,
env,
}
}
pub fn validate(&self) -> ModelResult<()> {
if self.image.trim().is_empty() {
return Err(ModelError::Invalid(
"container image cannot be empty".into(),
));
}
Ok(())
}
}
#[cfg(test)]
mod tests {
use super::*;
use std::path::PathBuf;
#[test]
fn built_in_validation_accepts_valid_values_and_rejects_empty_fields() {
TaskWorkload::Container(ContainerSpec {
image: "nginx:latest".into(),
command: None,
args: vec![],
env: Default::default(),
})
.validate()
.unwrap();
TaskWorkload::Embedded(EmbeddedSpec::new("v1").unwrap())
.validate()
.unwrap();
for image in ["", " \t"] {
let workload = TaskWorkload::Container(ContainerSpec {
image: image.into(),
command: None,
args: vec![],
env: Default::default(),
});
let error = workload.validate().unwrap_err();
assert!(error.to_string().contains("container image"));
}
let workload = TaskWorkload::Wasm(WasmSpec {
module: PathBuf::new(),
args: vec![],
env: Default::default(),
});
let error = workload.validate().unwrap_err();
assert!(error.to_string().contains("wasm module"));
assert!(EmbeddedSpec::new(" ").is_err());
}
#[test]
fn built_in_envelope_has_stable_gvk_and_rejects_unknown_fields() {
let workload = TaskWorkload::Embedded(EmbeddedSpec::new("build-42").unwrap());
let json = serde_json::to_value(&workload).unwrap();
assert_eq!(json["apiVersion"], "solti.io/v1");
assert_eq!(json["kind"], "Embedded");
assert_eq!(json["spec"], serde_json::json!({"revision": "build-42"}));
let mut envelope = serde_json::to_value(&workload).unwrap();
envelope["unexpected"] = serde_json::json!(true);
assert!(serde_json::from_value::<TaskWorkload>(envelope).is_err());
let mut spec = serde_json::to_value(workload).unwrap();
spec["spec"]["unexpected"] = serde_json::json!(true);
assert!(serde_json::from_value::<TaskWorkload>(spec).is_err());
}
#[test]
fn extension_roundtrips_gvk_and_application_owned_fields() {
let workload = TaskWorkload::Extension(
ExtensionWorkload::new(
"tasks.example.io/v1alpha1",
"ImageResize",
serde_json::json!({
"width": 1280,
"format": "webp",
"unexpectedToSolti": true,
"nested": { "applicationField": [1, 2, 3] }
}),
)
.unwrap(),
);
let json = serde_json::to_string(&workload).unwrap();
let back: TaskWorkload = serde_json::from_str(&json).unwrap();
assert_eq!(back, workload);
assert_eq!(back.api_version(), "tasks.example.io/v1alpha1");
assert_eq!(back.kind(), "ImageResize");
let extension = ExtensionWorkload::new(
"tasks.example.io/v1",
"Report",
serde_json::json!({ "format": "json" }),
)
.unwrap();
let json = serde_json::to_string(&extension).unwrap();
assert_eq!(
serde_json::from_str::<ExtensionWorkload>(&json).unwrap(),
extension
);
}
#[test]
fn extension_workload_rejects_reserved_solti_api_group() {
for api_version in [WORKLOAD_API_VERSION, "solti.io/v2"] {
let error =
ExtensionWorkload::new(api_version, "Custom", serde_json::json!({})).unwrap_err();
assert!(
error.to_string().contains("reserved"),
"apiVersion={api_version}"
);
}
}
#[test]
fn extension_workload_allows_builtin_kind_in_another_api_version() {
for kind in ["Subprocess", "Wasm", "Container", "Embedded"] {
let workload = TaskWorkload::Extension(
ExtensionWorkload::new(
"tasks.example.io/v1",
kind,
serde_json::json!({ "custom": true }),
)
.unwrap(),
);
let json = serde_json::to_string(&workload).unwrap();
let back: TaskWorkload = serde_json::from_str(&json).unwrap();
assert_eq!(back, workload, "kind={kind}");
}
}
#[test]
fn workload_gvk_uses_kubernetes_crd_validation() {
WorkloadTypeMeta::new("tasks.example.io/v1alpha1", "ImageResize").unwrap();
ExtensionWorkload::new("tasks.example.io/v1", "custom-kind", serde_json::json!({}))
.unwrap();
for api_version in [
"",
" solti.io/v1",
"bad/version/extra",
"example/v1",
"tasks.example.io/1v",
] {
assert!(
ExtensionWorkload::new(api_version, "Example", serde_json::json!({})).is_err(),
"apiVersion={api_version}"
);
}
for kind in ["", "1Example", "Bad Kind", "_Example"] {
assert!(
ExtensionWorkload::new("tasks.example.io/v1", kind, serde_json::json!({})).is_err(),
"kind={kind}"
);
}
}
#[test]
fn extension_workload_requires_object_spec() {
let error =
ExtensionWorkload::new("example.io/v1", "Example", serde_json::json!(42)).unwrap_err();
assert!(error.to_string().contains("JSON object"));
}
#[test]
fn constructors_build_specs_with_expected_fields() {
use crate::{Flag, SubprocessMode, TaskEnv};
let sub = SubprocessSpec::new(
SubprocessMode::Command {
command: "ls".into(),
args: vec!["-l".into()],
},
TaskEnv::default(),
Some(PathBuf::from("/tmp")),
Flag::enabled(),
);
assert!(matches!(sub.mode, SubprocessMode::Command { .. }));
assert_eq!(sub.cwd, Some(PathBuf::from("/tmp")));
let wasm = WasmSpec::new(
PathBuf::from("/m.wasm"),
vec!["--x".into()],
TaskEnv::default(),
);
assert_eq!(wasm.module, PathBuf::from("/m.wasm"));
assert_eq!(wasm.args, vec!["--x".to_string()]);
let cont = ContainerSpec::new(
"img:1".into(),
Some(vec!["sh".into()]),
vec!["-c".into()],
TaskEnv::default(),
);
assert_eq!(cont.image, "img:1");
assert_eq!(cont.command, Some(vec!["sh".to_string()]));
}
}
#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(rename_all = "camelCase", deny_unknown_fields)]
#[non_exhaustive]
pub struct SubprocessSpec {
pub mode: SubprocessMode,
#[serde(default, skip_serializing_if = "TaskEnv::is_empty")]
pub env: TaskEnv,
#[serde(skip_serializing_if = "Option::is_none")]
pub cwd: Option<PathBuf>,
#[serde(default)]
pub fail_on_non_zero: Flag,
}
impl SubprocessSpec {
pub fn new(
mode: SubprocessMode,
env: TaskEnv,
cwd: Option<PathBuf>,
fail_on_non_zero: Flag,
) -> Self {
Self {
mode,
env,
cwd,
fail_on_non_zero,
}
}
}
#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(rename_all = "camelCase", deny_unknown_fields)]
#[non_exhaustive]
pub struct WasmSpec {
#[cfg_attr(
feature = "schema",
schemars(schema_with = "crate::schema::path_string")
)]
pub module: PathBuf,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub args: Vec<String>,
#[serde(default, skip_serializing_if = "TaskEnv::is_empty")]
pub env: TaskEnv,
}
#[derive(Clone, Debug, Eq, PartialEq, Serialize, Deserialize)]
#[cfg_attr(feature = "schema", derive(schemars::JsonSchema))]
#[serde(rename_all = "camelCase", deny_unknown_fields)]
#[non_exhaustive]
pub struct ContainerSpec {
#[cfg_attr(
feature = "schema",
schemars(schema_with = "crate::schema::non_empty_string")
)]
pub image: String,
#[serde(skip_serializing_if = "Option::is_none")]
pub command: Option<Vec<String>>,
#[serde(default, skip_serializing_if = "Vec::is_empty")]
pub args: Vec<String>,
#[serde(default, skip_serializing_if = "TaskEnv::is_empty")]
pub env: TaskEnv,
}