use std::collections::BTreeMap;
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::workflow::EntryView;
use crate::workflow::PackageCandidate;
use crate::workflow::WORKFLOWS_CONFIG_KEY;
use crate::workflow::Workflow;
use crate::workflow::WorkflowId;
use crate::workflow::WorkflowIntent;
use crate::workflow::WorkflowRules;
use crate::workflow::WorkflowsConfig;
use crate::workflow::validate_package;
use quilt_uri::Host;
use quilt_uri::S3Uri;
async fn resolve_schema_url<R: Remote>(
remote: &R,
host: Option<&Host>,
config: &WorkflowsConfig,
workflow_id: &str,
schema_id: &str,
) -> Res<S3Uri> {
let declared = config.declared_schema_url(workflow_id, schema_id)?;
remote.resolve_url(host, &declared).await
}
async fn resolve_declared_schemas<R: Remote>(
remote: &R,
host: Option<&Host>,
config: &WorkflowsConfig,
workflow_id: &str,
) -> Res<BTreeMap<String, S3Uri>> {
let mut schemas = BTreeMap::new();
for key in ["metadata_schema", "entries_schema"] {
if let Some(schema_id) = config.schema_id(workflow_id, key)? {
let url = resolve_schema_url(remote, host, config, workflow_id, &schema_id).await?;
schemas.insert(schema_id, url);
}
}
Ok(schemas)
}
async fn resolve_named<R: Remote>(
remote: &R,
host: Option<&Host>,
config: &WorkflowsConfig,
config_uri: S3Uri,
id: String,
) -> Res<Option<Workflow>> {
let schemas = resolve_declared_schemas(remote, host, config, &id).await?;
Ok(Some(Workflow {
config: config_uri,
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 = resolve_schema_url(remote, host, config, 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.has_workflow(id) {
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)) => {
resolve_named(remote, host, parsed, 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 | WorkflowIntent::BucketDefault) => Ok(None),
(Some(parsed), WorkflowIntent::BucketDefault) => match parsed.bucket_default_id()? {
None => Ok(Some(Workflow { config, id: None })),
Some(id) => resolve_named(remote, host, parsed, config, id).await,
},
}
}
#[cfg(test)]
mod tests {
use super::*;
use crate::io::remote::mocks::MockRemote;
use crate::workflow::WorkflowInfo;
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.as_ref(),
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.as_ref(), WorkflowIntent::NoWorkflow, &uri).await?;
assert!(result.is_none());
let err = resolve_workflow(
&remote,
host.as_ref(),
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.as_ref(),
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.as_ref(), WorkflowIntent::NoWorkflow, &uri)
.await?
.unwrap();
assert_eq!(result.config, uri);
assert!(result.id.is_none());
let err = resolve_workflow(
&remote,
host.as_ref(),
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.as_ref(),
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.as_ref(), 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.as_ref(), 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.as_ref(), 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.as_ref(),
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.as_ref(),
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.as_ref(), 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.as_ref(), 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.as_ref(), 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_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.as_ref(), "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.as_ref(), "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.as_ref(), &uri).await?;
let parsed = parsed.expect("config present");
let rules = fetch_workflow_rules(&remote, host.as_ref(), &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.as_ref(), &uri).await?;
let parsed = parsed.expect("config present");
let rules = fetch_workflow_rules(&remote, host.as_ref(), &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.as_ref(), &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::object_hash::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");
}
}