use std::collections::BTreeMap;
use std::sync::LazyLock;
use jsonschema::Validator;
use serde::Deserialize;
use serde::Serialize;
use serde_json::Value;
use serde_yaml::Value as YamlValue;
use tokio::io::AsyncReadExt;
use crate::Error;
use crate::Res;
use crate::error::RemoteCatalogError;
use crate::io::remote::Remote;
use crate::manifest::ManifestRow;
use crate::manifest::Workflow;
use crate::manifest::WorkflowId;
use crate::workflow::EntryView;
use crate::workflow::PackageCandidate;
use crate::workflow::WorkflowRules;
use crate::workflow::compile_config_schema;
use crate::workflow::validate_package;
use quilt_uri::Host;
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) -> Res<()> {
use std::fmt::Write;
let document: Value = serde_json::to_value(yaml).map_err(|err| {
RemoteCatalogError::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(RemoteCatalogError::InvalidWorkflowsConfig(format!(
"workflows/config.yml does not satisfy the workflows config schema:{message}"
))
.into())
}
}
#[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) -> Res<WorkflowsConfig> {
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)
}
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)
}
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)
}
fn schema_id(&self, workflow_id: &str, key: &str) -> Res<Option<String>> {
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(Error::RemoteCatalog(RemoteCatalogError::Workflow(format!(
"`{key}` for workflow ID {workflow_id} must be a string"
)))),
},
_ => Err(Error::RemoteCatalog(RemoteCatalogError::Workflow(format!(
"Workflow {workflow_id} not found in workflows/config.yaml"
)))),
},
_ => Err(Error::RemoteCatalog(RemoteCatalogError::Workflow(
"Workflows not found in workflows/config.yaml".to_string(),
))),
}
}
fn declared_schema_url(&self, workflow_id: &str, schema_id: &str) -> Res<S3Uri> {
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(Error::RemoteCatalog(RemoteCatalogError::Workflow(format!(
"Schema {schema_id} doesn't have URL"
)))),
},
_ => Err(Error::RemoteCatalog(RemoteCatalogError::Workflow(format!(
"Schema {schema_id}, referenced by workflow {workflow_id} not found in workflows/config.yaml",
)))),
},
_ => Err(Error::RemoteCatalog(RemoteCatalogError::Workflow(
"Schemas not found in workflows/config.yaml".to_string(),
))),
}
}
async fn resolve_schema_url<R: Remote>(
&self,
remote: &R,
host: &Option<Host>,
workflow_id: &str,
schema_id: &str,
) -> Res<S3Uri> {
let declared = self.declared_schema_url(workflow_id, schema_id)?;
remote.resolve_url(host, &declared).await
}
fn declared_schema_uri(&self, workflow_id: &str, key: &str) -> Res<Option<S3Uri>> {
match self.schema_id(workflow_id, key)? {
Some(schema_id) => Ok(Some(self.declared_schema_url(workflow_id, &schema_id)?)),
None => Ok(None),
}
}
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(),
}
}
async fn resolve_declared_schemas<R: Remote>(
&self,
remote: &R,
host: &Option<Host>,
workflow_id: &str,
) -> Res<BTreeMap<String, S3Uri>> {
let mut schemas = BTreeMap::new();
for key in ["metadata_schema", "entries_schema"] {
if let Some(schema_id) = self.schema_id(workflow_id, key)? {
let url = self
.resolve_schema_url(remote, host, workflow_id, &schema_id)
.await?;
schemas.insert(schema_id, url);
}
}
Ok(schemas)
}
fn bucket_default_id(&self) -> Res<Option<String>> {
match self.raw.get("default_workflow") {
None => Ok(None),
Some(YamlValue::String(id)) => Ok(Some(id.clone())),
Some(_) => Err(Error::RemoteCatalog(RemoteCatalogError::Workflow(
"`default_workflow` in workflows/config.yaml must be a string".to_string(),
))),
}
}
async fn resolve_named<R: Remote>(
&self,
remote: &R,
host: &Option<Host>,
config: S3Uri,
id: String,
) -> Res<Option<Workflow>> {
let schemas = self.resolve_declared_schemas(remote, host, &id).await?;
Ok(Some(Workflow {
config,
id: Some(WorkflowId { id, schemas }),
}))
}
}
pub(crate) async fn fetch_workflows_config<R: Remote>(
remote: &R,
host: &Option<Host>,
uri: &S3Uri,
) -> Res<(S3Uri, Option<WorkflowsConfig>)> {
if !remote.exists(host, uri).await? {
return Ok((uri.clone(), None));
}
match remote.get_object_stream(host, uri).await {
Ok(stream) => {
let mut bytes = Vec::new();
stream
.body
.into_async_read()
.read_to_end(&mut bytes)
.await?;
let config = serde_yaml::from_slice::<Option<YamlValue>>(&bytes)?
.map(|yaml| WorkflowsConfig::from_yaml(&yaml))
.transpose()?;
Ok((stream.uri, config))
}
Err(err) => Err(err),
}
}
pub async fn fetch_workflows_config_for_bucket<R: Remote>(
remote: &R,
host: &Option<Host>,
bucket: &str,
) -> Res<Option<WorkflowsConfig>> {
let uri = S3Uri {
key: WORKFLOWS_CONFIG_KEY.to_string(),
bucket: bucket.to_string(),
version: None,
};
let (_, config) = fetch_workflows_config(remote, host, &uri).await?;
Ok(config)
}
async fn fetch_schema_doc<R: Remote>(remote: &R, host: &Option<Host>, uri: &S3Uri) -> Res<Value> {
let stream = remote.get_object_stream(host, uri).await?;
let mut bytes = Vec::new();
stream
.body
.into_async_read()
.read_to_end(&mut bytes)
.await?;
Ok(serde_json::from_slice(&bytes)?)
}
pub async fn fetch_workflow_rules<R: Remote>(
remote: &R,
host: &Option<Host>,
config: &WorkflowsConfig,
workflow_id: &str,
) -> Res<WorkflowRules> {
let metadata_schema =
fetch_schema_for_key(remote, host, config, workflow_id, "metadata_schema").await?;
let entries_schema =
fetch_schema_for_key(remote, host, config, workflow_id, "entries_schema").await?;
Ok(WorkflowRules {
handle_pattern: config.handle_pattern(workflow_id),
is_message_required: config.is_message_required(workflow_id),
metadata_schema,
entries_schema,
})
}
async fn fetch_schema_for_key<R: Remote>(
remote: &R,
host: &Option<Host>,
config: &WorkflowsConfig,
workflow_id: &str,
key: &str,
) -> Res<Option<Value>> {
match config.schema_id(workflow_id, key)? {
Some(schema_id) => {
let uri = config
.resolve_schema_url(remote, host, workflow_id, &schema_id)
.await?;
Ok(Some(fetch_schema_doc(remote, host, &uri).await?))
}
None => Ok(None),
}
}
pub(crate) fn entry_view(row: &ManifestRow) -> EntryView<'_> {
EntryView {
logical_key: row.logical_key.to_string_lossy(),
size: row.size,
meta: row.meta.as_ref().and_then(|meta| meta.get("user_meta")),
}
}
pub(crate) async fn validate_workflow<R: Remote>(
remote: &R,
host: &Option<Host>,
name: &str,
message: Option<&str>,
user_meta: Option<&Value>,
workflow: Option<&Workflow>,
entries: &[EntryView<'_>],
) -> Res<()> {
let config = match workflow {
Some(workflow) => {
fetch_workflows_config(remote, host, &workflow.config)
.await?
.1
}
None => None,
};
validate_workflow_with_config(
remote,
host,
name,
message,
user_meta,
workflow,
config.as_ref(),
entries,
)
.await
}
#[allow(clippy::too_many_arguments)]
pub(crate) async fn validate_workflow_with_config<R: Remote>(
remote: &R,
host: &Option<Host>,
name: &str,
message: Option<&str>,
user_meta: Option<&Value>,
workflow: Option<&Workflow>,
config: Option<&WorkflowsConfig>,
entries: &[EntryView<'_>],
) -> Res<()> {
let Some(workflow) = workflow else {
return Ok(());
};
let Some(config) = config else {
return Ok(());
};
let rules = match &workflow.id {
Some(workflow_id) => {
Some(fetch_workflow_rules(remote, host, config, &workflow_id.id).await?)
}
None => None,
};
let candidate = PackageCandidate {
name,
message,
user_meta,
entries,
};
validate_package(rules.as_ref(), config.is_workflow_required, &candidate)?;
Ok(())
}
#[allow(clippy::too_many_arguments)]
pub(crate) async fn validate_workflow_against_current_config<R: Remote>(
remote: &R,
host: &Option<Host>,
bucket: &str,
name: &str,
message: Option<&str>,
user_meta: Option<&Value>,
header_workflow: Option<&Workflow>,
entries: &[EntryView<'_>],
) -> Res<()> {
let workflow_id = header_workflow
.and_then(|workflow| workflow.id.as_ref())
.map(|id| id.id.as_str());
let config = fetch_workflows_config_for_bucket(remote, host, bucket).await?;
let Some(config) = config else {
return match workflow_id {
None => Ok(()),
Some(id) => Err(Error::RemoteCatalog(RemoteCatalogError::Workflow(format!(
"\"{id}\" workflow is specified, but no workflows config exists"
)))),
};
};
let rules = match workflow_id {
Some(id) => {
if config.workflow_entry(id).is_none() {
return Err(Error::RemoteCatalog(RemoteCatalogError::Workflow(format!(
"There is no \"{id}\" workflow in the config"
))));
}
Some(fetch_workflow_rules(remote, host, &config, id).await?)
}
None => None,
};
let candidate = PackageCandidate {
name,
message,
user_meta,
entries,
};
validate_package(rules.as_ref(), config.is_workflow_required, &candidate)?;
Ok(())
}
pub async fn resolve_workflow<R: Remote>(
remote: &R,
host: &Option<Host>,
intent: WorkflowIntent,
uri: &S3Uri,
) -> Res<Option<Workflow>> {
let (config, parsed) = fetch_workflows_config(remote, host, uri).await?;
resolve_workflow_from_config(remote, host, intent, config, parsed.as_ref()).await
}
pub(crate) async fn resolve_workflow_from_config<R: Remote>(
remote: &R,
host: &Option<Host>,
intent: WorkflowIntent,
config: S3Uri,
parsed: Option<&WorkflowsConfig>,
) -> Res<Option<Workflow>> {
match (parsed, intent) {
(Some(parsed), WorkflowIntent::Named(id)) => {
parsed.resolve_named(remote, host, config, id).await
}
(None, WorkflowIntent::Named(id)) => {
Err(Error::RemoteCatalog(RemoteCatalogError::Workflow(format!(
"There is no workflows config, but the workflow \"{id}\" is set"
))))
}
(Some(_), WorkflowIntent::NoWorkflow) => Ok(Some(Workflow { config, id: None })),
(None, WorkflowIntent::NoWorkflow) => Ok(None),
(None, WorkflowIntent::BucketDefault) => Ok(None),
(Some(parsed), WorkflowIntent::BucketDefault) => match parsed.bucket_default_id()? {
None => Ok(Some(Workflow { config, id: None })),
Some(id) => parsed.resolve_named(remote, host, config, id).await,
},
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::io::remote::mocks::MockRemote;
use test_log::test;
#[test(tokio::test)]
async fn test_missing_schemas_section() -> Res<()> {
let remote = MockRemote::default();
let host = None;
let uri: S3Uri = "s3://any/.quilt/workflows/config.yml".parse()?;
let config = r"
version: '1'
workflows:
foo:
name: Foo
metadata_schema: bar
";
remote
.put_object(&None, &uri, config.as_bytes().to_vec())
.await?;
let err = resolve_workflow(
&remote,
&host,
WorkflowIntent::Named("foo".to_string()),
&uri,
)
.await
.unwrap_err();
assert!(matches!(
err,
Error::RemoteCatalog(RemoteCatalogError::Workflow(_))
));
assert!(err.to_string().contains("Schemas not found"));
Ok(())
}
#[test(tokio::test)]
async fn test_no_config_yaml() -> Res<()> {
let remote = MockRemote::default();
let host = None;
let uri: S3Uri = "s3://any/.quilt/workflows/config.yml".parse()?;
let result = resolve_workflow(&remote, &host, WorkflowIntent::NoWorkflow, &uri).await?;
assert!(result.is_none());
let err = resolve_workflow(
&remote,
&host,
WorkflowIntent::Named("test-workflow".to_string()),
&uri,
)
.await
.unwrap_err();
assert!(matches!(
err,
Error::RemoteCatalog(RemoteCatalogError::Workflow(_))
));
assert!(err.to_string().contains("There is no workflows config"));
Ok(())
}
#[test(tokio::test)]
async fn test_with_config_yaml() -> Res<()> {
let remote = MockRemote::default();
let host = None;
let uri: S3Uri = "s3://any/.quilt/workflows/config.yml".parse()?;
let config = r"
version: '1'
workflows:
foo:
name: Foo
metadata_schema: bar
schemas:
bar:
url: s3://test-bucket/schemas/test.json
";
let schema_uri: S3Uri = "s3://test-bucket/schemas/test.json".parse()?;
let schema = b"{}";
remote
.put_object(&None, &uri, config.as_bytes().to_vec())
.await?;
remote
.put_object(&None, &schema_uri, schema.to_vec())
.await?;
let result = resolve_workflow(
&remote,
&host,
WorkflowIntent::Named("foo".to_string()),
&uri,
)
.await?
.unwrap();
assert_eq!(result.config, uri);
assert_eq!(
result.id.unwrap(),
WorkflowId {
id: "foo".to_string(),
schemas: BTreeMap::from([(
"bar".to_string(),
"s3://test-bucket/schemas/test.json".parse()?
)])
}
);
let result = resolve_workflow(&remote, &host, WorkflowIntent::NoWorkflow, &uri)
.await?
.unwrap();
assert_eq!(result.config, uri);
assert!(result.id.is_none());
let err = resolve_workflow(
&remote,
&host,
WorkflowIntent::Named("non-existent".to_string()),
&uri,
)
.await
.unwrap_err();
assert!(matches!(
err,
Error::RemoteCatalog(RemoteCatalogError::Workflow(_))
));
assert!(err.to_string().contains("Workflow non-existent not found"));
let err = resolve_workflow(&remote, &host, WorkflowIntent::Named(String::new()), &uri)
.await
.unwrap_err();
assert!(matches!(
err,
Error::RemoteCatalog(RemoteCatalogError::Workflow(_))
));
assert!(err.to_string().contains("Workflow not found"));
Ok(())
}
#[test(tokio::test)]
async fn test_bucket_default_no_config() -> Res<()> {
let remote = MockRemote::default();
let host = None;
let uri: S3Uri = "s3://any/.quilt/workflows/config.yml".parse()?;
let result = resolve_workflow(&remote, &host, WorkflowIntent::BucketDefault, &uri).await?;
assert!(result.is_none());
Ok(())
}
#[test(tokio::test)]
async fn test_bucket_default_without_key() -> Res<()> {
let remote = MockRemote::default();
let host = None;
let uri: S3Uri = "s3://any/.quilt/workflows/config.yml".parse()?;
let config = r"
version: '1'
workflows:
foo:
name: Foo
metadata_schema: bar
";
remote
.put_object(&None, &uri, config.as_bytes().to_vec())
.await?;
let result = resolve_workflow(&remote, &host, WorkflowIntent::BucketDefault, &uri)
.await?
.unwrap();
assert_eq!(result.config, uri);
assert!(result.id.is_none());
Ok(())
}
#[test(tokio::test)]
async fn test_bucket_default_with_valid_key() -> Res<()> {
let remote = MockRemote::default();
let host = None;
let uri: S3Uri = "s3://any/.quilt/workflows/config.yml".parse()?;
let config = r"
version: '1'
default_workflow: foo
workflows:
foo:
name: Foo
metadata_schema: bar
schemas:
bar:
url: s3://test-bucket/schemas/test.json
";
let schema_uri: S3Uri = "s3://test-bucket/schemas/test.json".parse()?;
remote
.put_object(&None, &uri, config.as_bytes().to_vec())
.await?;
remote
.put_object(&None, &schema_uri, b"{}".to_vec())
.await?;
let result = resolve_workflow(&remote, &host, WorkflowIntent::BucketDefault, &uri)
.await?
.unwrap();
assert_eq!(result.config, uri);
assert_eq!(
result.id.unwrap(),
WorkflowId {
id: "foo".to_string(),
schemas: BTreeMap::from([(
"bar".to_string(),
"s3://test-bucket/schemas/test.json".parse()?
)])
}
);
Ok(())
}
#[test(tokio::test)]
async fn test_resolve_stamps_both_declared_schemas() -> Res<()> {
let remote = MockRemote::default();
let host = None;
let uri: S3Uri = "s3://any/.quilt/workflows/config.yml".parse()?;
let config = r"
version: '1'
workflows:
dual:
name: Dual
metadata_schema: meta-id
entries_schema: entries-id
schemas:
meta-id:
url: s3://test-bucket/schemas/meta.json
entries-id:
url: s3://test-bucket/schemas/entries.json
";
remote
.put_object(&None, &uri, config.as_bytes().to_vec())
.await?;
for key in ["meta.json", "entries.json"] {
remote
.put_object(
&None,
&format!("s3://test-bucket/schemas/{key}").parse()?,
b"{}".to_vec(),
)
.await?;
}
let result = resolve_workflow(
&remote,
&host,
WorkflowIntent::Named("dual".to_string()),
&uri,
)
.await?
.unwrap();
assert_eq!(
result.id.unwrap(),
WorkflowId {
id: "dual".to_string(),
schemas: BTreeMap::from([
(
"meta-id".to_string(),
"s3://test-bucket/schemas/meta.json".parse()?
),
(
"entries-id".to_string(),
"s3://test-bucket/schemas/entries.json".parse()?
),
]),
}
);
Ok(())
}
#[test(tokio::test)]
async fn test_resolve_stamps_entries_only_schema() -> Res<()> {
let remote = MockRemote::default();
let host = None;
let uri: S3Uri = "s3://any/.quilt/workflows/config.yml".parse()?;
let config = r"
version: '1'
workflows:
entries-only:
name: Entries only
entries_schema: entries-id
schemas:
entries-id:
url: s3://test-bucket/schemas/entries.json
";
remote
.put_object(&None, &uri, config.as_bytes().to_vec())
.await?;
remote
.put_object(
&None,
&"s3://test-bucket/schemas/entries.json".parse()?,
b"{}".to_vec(),
)
.await?;
let result = resolve_workflow(
&remote,
&host,
WorkflowIntent::Named("entries-only".to_string()),
&uri,
)
.await?
.unwrap();
assert_eq!(
result.id.unwrap(),
WorkflowId {
id: "entries-only".to_string(),
schemas: BTreeMap::from([(
"entries-id".to_string(),
"s3://test-bucket/schemas/entries.json".parse()?
)]),
}
);
Ok(())
}
#[test(tokio::test)]
async fn test_bucket_default_missing_referenced_workflow() -> Res<()> {
let remote = MockRemote::default();
let host = None;
let uri: S3Uri = "s3://any/.quilt/workflows/config.yml".parse()?;
let config = r"
version: '1'
default_workflow: ghost
workflows:
foo:
name: Foo
metadata_schema: bar
";
remote
.put_object(&None, &uri, config.as_bytes().to_vec())
.await?;
let err = resolve_workflow(&remote, &host, WorkflowIntent::BucketDefault, &uri)
.await
.unwrap_err();
assert!(matches!(
err,
Error::RemoteCatalog(RemoteCatalogError::Workflow(_))
));
assert!(err.to_string().contains("Workflow ghost not found"));
Ok(())
}
#[test(tokio::test)]
async fn test_bucket_default_non_string_key() -> Res<()> {
let remote = MockRemote::default();
let host = None;
let uri: S3Uri = "s3://any/.quilt/workflows/config.yml".parse()?;
let config = r"
version: '1'
default_workflow: [not, a, string]
workflows:
foo:
name: Foo
metadata_schema: bar
";
remote
.put_object(&None, &uri, config.as_bytes().to_vec())
.await?;
let err = resolve_workflow(&remote, &host, WorkflowIntent::BucketDefault, &uri)
.await
.unwrap_err();
assert!(matches!(
err,
Error::RemoteCatalog(RemoteCatalogError::InvalidWorkflowsConfig(_))
));
let message = err.to_string();
assert!(
message.contains("does not satisfy the workflows config schema"),
"unexpected message: {message}"
);
assert!(
message.contains("/default_workflow"),
"violation must name the offending path, got: {message}"
);
Ok(())
}
#[test(tokio::test)]
async fn test_bucket_default_explicit_null_key() -> Res<()> {
let remote = MockRemote::default();
let host = None;
let uri: S3Uri = "s3://any/.quilt/workflows/config.yml".parse()?;
let config = r"
version: '1'
default_workflow:
workflows:
foo:
name: Foo
metadata_schema: bar
";
remote
.put_object(&None, &uri, config.as_bytes().to_vec())
.await?;
let err = resolve_workflow(&remote, &host, WorkflowIntent::BucketDefault, &uri)
.await
.unwrap_err();
assert!(matches!(
err,
Error::RemoteCatalog(RemoteCatalogError::InvalidWorkflowsConfig(_))
));
let message = err.to_string();
assert!(
message.contains("does not satisfy the workflows config schema"),
"unexpected message: {message}"
);
assert!(
message.contains("/default_workflow"),
"violation must name the offending path, got: {message}"
);
Ok(())
}
#[test(tokio::test)]
async 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,
Error::RemoteCatalog(RemoteCatalogError::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(tokio::test)]
async fn test_fetch_workflows_config_for_bucket_present() -> Res<()> {
let remote = MockRemote::default();
let host = None;
let uri: S3Uri = "s3://my-bucket/.quilt/workflows/config.yml".parse()?;
let config = r"
version: '1'
default_workflow: foo
workflows:
foo:
name: Foo
";
remote
.put_object(&None, &uri, config.as_bytes().to_vec())
.await?;
let parsed = fetch_workflows_config_for_bucket(&remote, &host, "my-bucket")
.await?
.expect("config present → Some");
assert_eq!(parsed.default_workflow, Some("foo".to_string()));
assert_eq!(
parsed.workflows,
vec![WorkflowInfo {
id: "foo".to_string(),
name: Some("Foo".to_string()),
description: None,
}]
);
Ok(())
}
#[test(tokio::test)]
async fn test_fetch_workflows_config_for_bucket_absent() -> Res<()> {
let remote = MockRemote::default();
let host = None;
let result = fetch_workflows_config_for_bucket(&remote, &host, "empty-bucket").await?;
assert!(result.is_none());
Ok(())
}
#[test(tokio::test)]
async fn test_fetch_workflow_rules() -> Res<()> {
let remote = MockRemote::default();
let host = None;
let uri: S3Uri = "s3://any/.quilt/workflows/config.yml".parse()?;
let config = r#"
version: "1"
workflows:
foo:
name: Foo
handle_pattern: "^team/"
is_message_required: true
metadata_schema: meta
entries_schema: entries
schemas:
meta:
url: s3://test-bucket/schemas/meta.json
entries:
url: s3://test-bucket/schemas/entries.json
"#;
let meta_uri: S3Uri = "s3://test-bucket/schemas/meta.json".parse()?;
let entries_uri: S3Uri = "s3://test-bucket/schemas/entries.json".parse()?;
remote
.put_object(&None, &uri, config.as_bytes().to_vec())
.await?;
remote
.put_object(
&None,
&meta_uri,
br#"{"type": "object", "required": ["owner"]}"#.to_vec(),
)
.await?;
remote
.put_object(&None, &entries_uri, br#"{"type": "array"}"#.to_vec())
.await?;
let (_, parsed) = fetch_workflows_config(&remote, &host, &uri).await?;
let parsed = parsed.expect("config present");
let rules = fetch_workflow_rules(&remote, &host, &parsed, "foo").await?;
assert_eq!(
rules,
WorkflowRules {
handle_pattern: Some("^team/".to_string()),
is_message_required: true,
metadata_schema: Some(serde_json::json!({
"type": "object", "required": ["owner"]
})),
entries_schema: Some(serde_json::json!({ "type": "array" })),
}
);
Ok(())
}
#[test(tokio::test)]
async fn test_fetch_workflow_rules_no_schemas() -> Res<()> {
let remote = MockRemote::default();
let host = None;
let uri: S3Uri = "s3://any/.quilt/workflows/config.yml".parse()?;
let config = r"
version: '1'
workflows:
bare:
name: Bare
";
remote
.put_object(&None, &uri, config.as_bytes().to_vec())
.await?;
let (_, parsed) = fetch_workflows_config(&remote, &host, &uri).await?;
let parsed = parsed.expect("config present");
let rules = fetch_workflow_rules(&remote, &host, &parsed, "bare").await?;
assert_eq!(
rules,
WorkflowRules {
handle_pattern: None,
is_message_required: false,
metadata_schema: None,
entries_schema: None,
}
);
Ok(())
}
#[test(tokio::test)]
async fn test_non_string_schema_key_rejected_by_config_schema() -> Res<()> {
let remote = MockRemote::default();
let host = None;
let uri: S3Uri = "s3://any/.quilt/workflows/config.yml".parse()?;
let config = r"
version: '1'
workflows:
foo:
name: Foo
metadata_schema: ~
";
remote
.put_object(&None, &uri, config.as_bytes().to_vec())
.await?;
let err = fetch_workflows_config(&remote, &host, &uri)
.await
.unwrap_err();
assert!(matches!(
err,
Error::RemoteCatalog(RemoteCatalogError::InvalidWorkflowsConfig(_))
));
let message = err.to_string();
assert!(
message.contains("does not satisfy the workflows config schema"),
"unexpected message: {message}"
);
assert!(
message.contains("/workflows/foo/metadata_schema"),
"violation must name the offending path, got: {message}"
);
Ok(())
}
#[cfg(unix)]
#[test]
fn test_entry_view_non_utf8_logical_key_is_lossy() {
use std::ffi::OsStr;
use std::os::unix::ffi::OsStrExt;
use std::path::PathBuf;
use crate::checksum::ObjectHash;
let row = ManifestRow {
logical_key: PathBuf::from(OsStr::from_bytes(b"bad-\xFF.txt")),
physical_key: "s3://bucket/bad".to_string(),
hash: ObjectHash::default(),
size: 1,
meta: None,
};
let view = entry_view(&row);
assert_eq!(view.logical_key, "bad-\u{FFFD}.txt");
}
#[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,
Error::RemoteCatalog(RemoteCatalogError::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,
Error::RemoteCatalog(RemoteCatalogError::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(())
}
}