use std::sync::LazyLock;
use jsonschema::Validator;
use serde::Deserialize;
use serde::Serialize;
use serde_json::Value;
use serde_yaml::Value as YamlValue;
use super::error::ConfigError;
use super::validate::compile_config_schema;
use quilt_uri::S3Uri;
const CONFIG_SCHEMA: &str = include_str!("config-1.schema.json");
pub const WORKFLOWS_CONFIG_KEY: &str = ".quilt/workflows/config.yml";
static CONFIG_VALIDATOR: LazyLock<Validator> = LazyLock::new(|| {
let schema: Value = serde_json::from_str(CONFIG_SCHEMA)
.expect("vendored workflows-config schema is valid JSON");
compile_config_schema(&schema)
});
fn validate_config_document(yaml: &YamlValue) -> Result<(), ConfigError> {
use std::fmt::Write;
let document: Value = serde_json::to_value(yaml).map_err(|err| {
ConfigError::InvalidWorkflowsConfig(format!(
"workflows/config.yml could not be converted for schema validation: {err}"
))
})?;
let mut message = String::new();
for err in CONFIG_VALIDATOR.iter_errors(&document) {
let _ = write!(message, "\n - {err} (at {})", err.instance_path());
}
if message.is_empty() {
Ok(())
} else {
Err(ConfigError::InvalidWorkflowsConfig(format!(
"workflows/config.yml does not satisfy the workflows config schema:{message}"
)))
}
}
#[derive(Debug, Clone, PartialEq, Serialize, Deserialize)]
#[serde(tag = "kind", content = "id", rename_all = "kebab-case")]
pub enum WorkflowIntent {
BucketDefault,
NoWorkflow,
Named(String),
}
impl WorkflowIntent {
pub fn from_optional_id(id: Option<&str>) -> Self {
match id.map(str::trim) {
Some(id) if !id.is_empty() => Self::Named(id.to_string()),
_ => Self::BucketDefault,
}
}
}
#[derive(Debug, Clone, PartialEq)]
pub struct WorkflowInfo {
pub id: String,
pub name: Option<String>,
pub description: Option<String>,
}
#[derive(Debug, Clone, PartialEq, Eq, Default)]
pub struct WorkflowSchemaUris {
pub metadata_schema: Option<S3Uri>,
pub entries_schema: Option<S3Uri>,
}
#[derive(Debug, Clone, PartialEq)]
pub struct WorkflowsConfig {
pub default_workflow: Option<String>,
pub is_workflow_required: bool,
pub workflows: Vec<WorkflowInfo>,
raw: YamlValue,
}
impl WorkflowsConfig {
pub fn from_yaml(yaml: &YamlValue) -> Result<WorkflowsConfig, ConfigError> {
validate_config_document(yaml)?;
let default_workflow = yaml
.get("default_workflow")
.and_then(YamlValue::as_str)
.map(String::from);
let is_workflow_required = yaml
.get("is_workflow_required")
.and_then(YamlValue::as_bool)
.unwrap_or(true);
let workflows = yaml
.get("workflows")
.and_then(YamlValue::as_mapping)
.map(|workflows| {
workflows
.iter()
.filter_map(|(id, entry)| {
Some(WorkflowInfo {
id: id.as_str()?.to_string(),
name: entry
.get("name")
.and_then(YamlValue::as_str)
.map(String::from),
description: entry
.get("description")
.and_then(YamlValue::as_str)
.map(String::from),
})
})
.collect()
})
.unwrap_or_default();
Ok(WorkflowsConfig {
default_workflow,
is_workflow_required,
workflows,
raw: yaml.clone(),
})
}
fn workflow_entry(&self, workflow_id: &str) -> Option<&YamlValue> {
self.raw.get("workflows")?.get(workflow_id)
}
pub(crate) fn has_workflow(&self, workflow_id: &str) -> bool {
self.workflow_entry(workflow_id).is_some()
}
pub(crate) fn handle_pattern(&self, workflow_id: &str) -> Option<String> {
self.workflow_entry(workflow_id)?
.get("handle_pattern")
.and_then(YamlValue::as_str)
.map(String::from)
}
pub(crate) fn is_message_required(&self, workflow_id: &str) -> bool {
self.workflow_entry(workflow_id)
.and_then(|workflow| workflow.get("is_message_required"))
.and_then(YamlValue::as_bool)
.unwrap_or(false)
}
pub(crate) fn schema_id(
&self,
workflow_id: &str,
key: &str,
) -> Result<Option<String>, ConfigError> {
match self.raw.get("workflows") {
Some(YamlValue::Mapping(workflows)) => match workflows.get(workflow_id) {
Some(YamlValue::Mapping(workflow)) => match workflow.get(key) {
Some(YamlValue::String(schema_id)) => Ok(Some(schema_id.clone())),
None => Ok(None),
Some(_) => Err(ConfigError::Workflow(format!(
"`{key}` for workflow ID {workflow_id} must be a string"
))),
},
_ => Err(ConfigError::Workflow(format!(
"Workflow {workflow_id} not found in workflows/config.yaml"
))),
},
_ => Err(ConfigError::Workflow(
"Workflows not found in workflows/config.yaml".to_string(),
)),
}
}
pub(crate) fn declared_schema_url(
&self,
workflow_id: &str,
schema_id: &str,
) -> Result<S3Uri, ConfigError> {
match self.raw.get("schemas") {
Some(YamlValue::Mapping(schemas)) => match schemas.get(schema_id) {
Some(YamlValue::Mapping(schema)) => match schema.get("url") {
Some(YamlValue::String(url)) => Ok(url.parse()?),
_ => Err(ConfigError::Workflow(format!(
"Schema {schema_id} doesn't have URL"
))),
},
_ => Err(ConfigError::Workflow(format!(
"Schema {schema_id}, referenced by workflow {workflow_id} not found in workflows/config.yaml",
))),
},
_ => Err(ConfigError::Workflow(
"Schemas not found in workflows/config.yaml".to_string(),
)),
}
}
fn declared_schema_uri(
&self,
workflow_id: &str,
key: &str,
) -> Result<Option<S3Uri>, ConfigError> {
match self.schema_id(workflow_id, key)? {
Some(schema_id) => Ok(Some(self.declared_schema_url(workflow_id, &schema_id)?)),
None => Ok(None),
}
}
#[must_use]
pub fn schema_uris(&self, workflow_id: &str) -> WorkflowSchemaUris {
WorkflowSchemaUris {
metadata_schema: self
.declared_schema_uri(workflow_id, "metadata_schema")
.ok()
.flatten(),
entries_schema: self
.declared_schema_uri(workflow_id, "entries_schema")
.ok()
.flatten(),
}
}
pub(crate) fn bucket_default_id(&self) -> Result<Option<String>, ConfigError> {
match self.raw.get("default_workflow") {
None => Ok(None),
Some(YamlValue::String(id)) => Ok(Some(id.clone())),
Some(_) => Err(ConfigError::Workflow(
"`default_workflow` in workflows/config.yaml must be a string".to_string(),
)),
}
}
}
#[cfg(test)]
mod tests {
use super::*;
use test_log::test;
type Res<T = ()> = Result<T, Box<dyn std::error::Error>>;
#[test]
fn test_non_string_mapping_key_is_invalid_config() -> Res<()> {
let yaml: YamlValue = serde_yaml::from_str(
r"
1: not-a-string-key
version: '1'
",
)?;
let err = WorkflowsConfig::from_yaml(&yaml).unwrap_err();
assert!(
matches!(err, ConfigError::InvalidWorkflowsConfig(_)),
"expected InvalidWorkflowsConfig, got: {err:?}"
);
Ok(())
}
#[test]
fn test_workflows_config_parse_rich() -> Res<()> {
let yaml: YamlValue = serde_yaml::from_str(
r#"
version: "1"
is_workflow_required: true
default_workflow: dummy
workflows:
dummy:
name: Dummy workflow
description: Do nothing.
alpha:
name: Alpha
description: First workflow.
metadata_schema: alpha-schema
schemas:
alpha-schema:
url: s3://sandbox/schemas/alpha.json
"#,
)?;
let config = WorkflowsConfig::from_yaml(&yaml)?;
assert_eq!(config.default_workflow, Some("dummy".to_string()));
assert!(config.is_workflow_required);
assert_eq!(
config.workflows,
vec![
WorkflowInfo {
id: "dummy".to_string(),
name: Some("Dummy workflow".to_string()),
description: Some("Do nothing.".to_string()),
},
WorkflowInfo {
id: "alpha".to_string(),
name: Some("Alpha".to_string()),
description: Some("First workflow.".to_string()),
},
]
);
Ok(())
}
#[test]
fn test_workflows_config_required_defaults_true() -> Res<()> {
let yaml: YamlValue = serde_yaml::from_str(
r"
version: '1'
workflows:
foo:
name: Foo
metadata_schema: bar
",
)?;
let config = WorkflowsConfig::from_yaml(&yaml)?;
assert!(config.is_workflow_required);
assert_eq!(config.default_workflow, None);
Ok(())
}
#[test]
fn test_workflows_config_required_explicit_false() -> Res<()> {
let yaml: YamlValue = serde_yaml::from_str(
r"
version: '1'
is_workflow_required: false
workflows:
foo:
name: Foo
",
)?;
let config = WorkflowsConfig::from_yaml(&yaml)?;
assert!(!config.is_workflow_required);
assert_eq!(
config.workflows,
vec![WorkflowInfo {
id: "foo".to_string(),
name: Some("Foo".to_string()),
description: None,
}]
);
Ok(())
}
#[test]
fn test_from_optional_id() {
assert_eq!(
WorkflowIntent::from_optional_id(None),
WorkflowIntent::BucketDefault
);
assert_eq!(
WorkflowIntent::from_optional_id(Some("x")),
WorkflowIntent::Named("x".to_string())
);
assert_eq!(
WorkflowIntent::from_optional_id(Some("")),
WorkflowIntent::BucketDefault
);
assert_eq!(
WorkflowIntent::from_optional_id(Some(" ")),
WorkflowIntent::BucketDefault
);
assert_eq!(
WorkflowIntent::from_optional_id(Some(" x ")),
WorkflowIntent::Named("x".to_string())
);
}
#[test]
fn test_workflow_intent_serde_round_trip() -> Res<()> {
for (intent, wire) in [
(
WorkflowIntent::BucketDefault,
serde_json::json!({ "kind": "bucket-default" }),
),
(
WorkflowIntent::NoWorkflow,
serde_json::json!({ "kind": "no-workflow" }),
),
(
WorkflowIntent::Named("foo".to_string()),
serde_json::json!({ "kind": "named", "id": "foo" }),
),
] {
assert_eq!(serde_json::to_value(&intent)?, wire);
assert_eq!(serde_json::from_value::<WorkflowIntent>(wire)?, intent);
}
Ok(())
}
#[test]
fn test_quoted_is_message_required_rejected_by_config_schema() {
let yaml: YamlValue = serde_yaml::from_str(
r#"
version: "1"
workflows:
foo:
name: Foo
is_message_required: "true"
"#,
)
.expect("valid YAML");
let err = WorkflowsConfig::from_yaml(&yaml).unwrap_err();
assert!(matches!(err, ConfigError::InvalidWorkflowsConfig(_)));
let message = err.to_string();
assert!(
message.contains("does not satisfy the workflows config schema"),
"unexpected message: {message}"
);
assert!(
message.contains("/workflows/foo/is_message_required"),
"violation must name the offending path, got: {message}"
);
}
#[test]
fn test_list_handle_pattern_rejected_by_config_schema() {
let yaml: YamlValue = serde_yaml::from_str(
r#"
version: "1"
workflows:
foo:
name: Foo
handle_pattern: ["^team/"]
"#,
)
.expect("valid YAML");
let err = WorkflowsConfig::from_yaml(&yaml).unwrap_err();
assert!(matches!(err, ConfigError::InvalidWorkflowsConfig(_)));
let message = err.to_string();
assert!(
message.contains("/workflows/foo/handle_pattern"),
"violation must name the offending path, got: {message}"
);
}
#[test]
fn test_schema_uris_both_one_none_and_unknown() -> Res<()> {
let yaml: YamlValue = serde_yaml::from_str(
r#"
version: "1"
workflows:
both:
name: Both
metadata_schema: meta
entries_schema: entries
meta-only:
name: Meta only
metadata_schema: meta
none:
name: None
schemas:
meta:
url: s3://schemas-bucket/meta.json
entries:
url: s3://schemas-bucket/entries.json
"#,
)?;
let config = WorkflowsConfig::from_yaml(&yaml)?;
assert_eq!(
config.schema_uris("both"),
WorkflowSchemaUris {
metadata_schema: Some("s3://schemas-bucket/meta.json".parse()?),
entries_schema: Some("s3://schemas-bucket/entries.json".parse()?),
}
);
assert_eq!(
config.schema_uris("meta-only"),
WorkflowSchemaUris {
metadata_schema: Some("s3://schemas-bucket/meta.json".parse()?),
entries_schema: None,
}
);
assert_eq!(config.schema_uris("none"), WorkflowSchemaUris::default());
assert_eq!(config.schema_uris("ghost"), WorkflowSchemaUris::default());
Ok(())
}
#[test]
fn test_schema_uris_resolve_independently() -> Res<()> {
let yaml: YamlValue = serde_yaml::from_str(
r#"
version: "1"
workflows:
partial:
name: Partial
metadata_schema: meta
entries_schema: ghost
schemas:
meta:
url: s3://schemas-bucket/meta.json
"#,
)?;
let config = WorkflowsConfig::from_yaml(&yaml)?;
assert_eq!(
config.schema_uris("partial"),
WorkflowSchemaUris {
metadata_schema: Some("s3://schemas-bucket/meta.json".parse()?),
entries_schema: None,
}
);
Ok(())
}
#[test]
fn test_schema_uris_missing_schemas_section_degrades_to_none() -> Res<()> {
let yaml: YamlValue = serde_yaml::from_str(
r"
version: '1'
workflows:
foo:
name: Foo
metadata_schema: bar
",
)?;
let config = WorkflowsConfig::from_yaml(&yaml)?;
assert_eq!(config.schema_uris("foo"), WorkflowSchemaUris::default());
Ok(())
}
#[test]
fn test_valid_config_with_format_annotations_parses() -> Res<()> {
let yaml: YamlValue = serde_yaml::from_str(
r#"
version: "1"
is_workflow_required: true
default_workflow: foo
workflows:
foo:
name: Foo
handle_pattern: "^team/"
is_message_required: true
metadata_schema: meta
schemas:
meta:
url: s3://bucket/schemas/meta.json
"#,
)?;
let config = WorkflowsConfig::from_yaml(&yaml)?;
assert_eq!(config.default_workflow, Some("foo".to_string()));
assert!(config.is_workflow_required);
assert!(config.is_message_required("foo"));
assert_eq!(config.handle_pattern("foo"), Some("^team/".to_string()));
Ok(())
}
}