use super::*;
use test_log::test;
use aws_sdk_s3::primitives::ByteStream;
use crate::io::remote::WorkflowIntent;
use crate::io::remote::mocks::MockRemote;
use crate::lineage::DomainLineageIo;
use crate::lineage::Home;
use crate::lineage::PackageLineageIo;
use crate::manifest::ManifestHeader;
use crate::object_hash::ObjectHash;
use crate::paths::DomainPaths;
use crate::workflow::RuleViolation;
use crate::workflow::WorkflowValidationError;
#[test(tokio::test)]
async fn test_set_remote_on_local_package() -> Res {
let (home, _temp_dir1) = Home::from_temp_dir()?;
let (paths, _temp_dir2) = DomainPaths::from_temp_dir()?;
let storage = LocalStorage::new();
let remote = MockRemote::default();
let namespace: Namespace = ("test", "local").into();
paths
.scaffold_for_installing(&storage, &home, &namespace)
.await?;
let lineage_json = r#"{
"packages": {
"test/local": {
"commit": null,
"remote": null,
"base_hash": "",
"latest_hash": "",
"paths": {}
}
},
"home": "/tmp/working_dir"
}"#;
storage
.write_byte_stream(&paths.lineage(), lineage_json.as_bytes().to_vec().into())
.await?;
let domain_lineage_io = DomainLineageIo::new(paths.lineage());
let package = InstalledPackage {
lineage: PackageLineageIo::new(domain_lineage_io, namespace.clone()),
paths,
remote,
storage,
namespace,
};
package
.set_remote(
"my-bucket".to_string(),
Some("example.com".parse()?),
WorkflowIntent::BucketDefault,
)
.await?;
let lineage = package.lineage().await?;
let remote_uri = lineage
.remote_uri
.as_ref()
.expect("remote_uri should be set");
assert_eq!(
remote_uri.origin.as_ref().unwrap().to_string(),
"example.com"
);
assert_eq!(remote_uri.bucket, "my-bucket");
assert_eq!(remote_uri.hash, "");
Ok(())
}
#[test(tokio::test)]
async fn test_set_remote_empty_bucket_error() -> Res {
let (home, _temp_dir1) = Home::from_temp_dir()?;
let (paths, _temp_dir2) = DomainPaths::from_temp_dir()?;
let storage = LocalStorage::new();
let remote = MockRemote::default();
let namespace: Namespace = ("test", "local").into();
paths
.scaffold_for_installing(&storage, &home, &namespace)
.await?;
let lineage_json = r#"{
"packages": {
"test/local": {
"commit": null,
"remote": null,
"base_hash": "",
"latest_hash": "",
"paths": {}
}
},
"home": "/tmp/working_dir"
}"#;
storage
.write_byte_stream(&paths.lineage(), lineage_json.as_bytes().to_vec().into())
.await?;
let domain_lineage_io = DomainLineageIo::new(paths.lineage());
let package = InstalledPackage {
lineage: PackageLineageIo::new(domain_lineage_io, namespace.clone()),
paths,
remote,
storage,
namespace,
};
let result = package
.set_remote(
String::new(),
Some("example.com".parse()?),
WorkflowIntent::BucketDefault,
)
.await;
assert!(result.is_err());
assert!(
result
.unwrap_err()
.to_string()
.contains("Bucket cannot be empty"),
"Error should mention empty bucket"
);
Ok(())
}
#[test(tokio::test)]
async fn test_set_remote_rejects_unreachable_bucket() -> Res {
use crate::error::RemoteCatalogError;
struct BadBucketRemote;
impl Remote for BadBucketRemote {
async fn exists(&self, _host: Option<&Host>, _s3_uri: &S3Uri) -> Res<bool> {
unreachable!("test only exercises verify_bucket")
}
async fn get_object_stream(
&self,
_host: Option<&Host>,
_s3_uri: &S3Uri,
) -> Res<crate::io::remote::RemoteObjectStream> {
unreachable!("test only exercises verify_bucket")
}
async fn resolve_url(&self, _host: Option<&Host>, _s3_uri: &S3Uri) -> Res<S3Uri> {
unreachable!("test only exercises verify_bucket")
}
async fn put_object(
&self,
_host: Option<&Host>,
_s3_uri: &S3Uri,
_contents: impl Into<aws_sdk_s3::primitives::ByteStream>,
) -> Res {
unreachable!("test only exercises verify_bucket")
}
async fn upload_file(
&self,
_host_config: &crate::io::remote::HostConfig,
_source_path: impl AsRef<std::path::Path>,
_dest_uri: &S3Uri,
_size: u64,
) -> Res<(S3Uri, ObjectHash)> {
unreachable!("test only exercises verify_bucket")
}
async fn host_config(&self, _host: Option<&Host>) -> Res<crate::io::remote::HostConfig> {
Ok(crate::io::remote::HostConfig::default())
}
async fn verify_bucket(&self, bucket: &str) -> Res {
Err(RemoteCatalogError::BucketUnreachable(bucket.to_string()).into())
}
}
let (home, _temp_dir1) = Home::from_temp_dir()?;
let (paths, _temp_dir2) = DomainPaths::from_temp_dir()?;
let storage = LocalStorage::new();
let namespace: Namespace = ("test", "badbucket").into();
paths
.scaffold_for_installing(&storage, &home, &namespace)
.await?;
let lineage_json = r#"{
"packages": {
"test/badbucket": {
"commit": null,
"remote": null,
"base_hash": "",
"latest_hash": "",
"paths": {}
}
},
"home": "/tmp/working_dir"
}"#;
storage
.write_byte_stream(&paths.lineage(), lineage_json.as_bytes().to_vec().into())
.await?;
let domain_lineage_io = DomainLineageIo::new(paths.lineage());
let package = InstalledPackage {
lineage: PackageLineageIo::new(domain_lineage_io, namespace.clone()),
paths,
remote: BadBucketRemote,
storage,
namespace,
};
let result = package
.set_remote(
"typo-bucket".to_string(),
Some("example.com".parse()?),
WorkflowIntent::BucketDefault,
)
.await;
assert!(result.is_err());
let msg = result.unwrap_err().to_string();
assert!(
msg.contains("typo-bucket") && msg.contains("not reachable"),
"error should name the bucket and say it's unreachable, got: {msg}"
);
let lineage = package.lineage().await?;
assert!(
lineage.remote_uri.is_none(),
"remote_uri should not be persisted when verify_bucket fails",
);
Ok(())
}
#[test(tokio::test)]
async fn test_set_remote_rejects_change_on_pushed_package() -> Res {
let (home, _temp_dir1) = Home::from_temp_dir()?;
let (paths, _temp_dir2) = DomainPaths::from_temp_dir()?;
let storage = LocalStorage::new();
let remote = MockRemote::default();
let namespace: Namespace = ("test", "overwrite").into();
paths
.scaffold_for_installing(&storage, &home, &namespace)
.await?;
let lineage_json = r#"{
"packages": {
"test/overwrite": {
"commit": null,
"remote": {
"bucket": "old-bucket",
"namespace": "test/overwrite",
"hash": "abc123",
"origin": "old.host"
},
"base_hash": "abc123",
"latest_hash": "abc123",
"paths": {}
}
},
"home": "/tmp/working_dir"
}"#;
storage
.write_byte_stream(&paths.lineage(), lineage_json.as_bytes().to_vec().into())
.await?;
let domain_lineage_io = DomainLineageIo::new(paths.lineage());
let package = InstalledPackage {
lineage: PackageLineageIo::new(domain_lineage_io, namespace.clone()),
paths,
remote,
storage,
namespace,
};
let result = package
.set_remote(
"new-bucket".to_string(),
Some("new.host".parse()?),
WorkflowIntent::BucketDefault,
)
.await;
assert!(result.is_err());
assert!(
result
.unwrap_err()
.to_string()
.contains("Cannot change remote"),
"Should reject changing remote on a pushed package"
);
Ok(())
}
#[test(tokio::test)]
async fn test_set_remote_is_idempotent_on_pushed_package() -> Res {
let (home, _temp_dir1) = Home::from_temp_dir()?;
let (paths, _temp_dir2) = DomainPaths::from_temp_dir()?;
let storage = LocalStorage::new();
let remote = MockRemote::default();
let namespace: Namespace = ("test", "idempotent").into();
paths
.scaffold_for_installing(&storage, &home, &namespace)
.await?;
let lineage_json = r#"{
"packages": {
"test/idempotent": {
"commit": null,
"remote": {
"bucket": "my-bucket",
"namespace": "test/idempotent",
"hash": "abc123",
"origin": "my.host"
},
"base_hash": "abc123",
"latest_hash": "abc123",
"paths": {}
}
},
"home": "/tmp/working_dir"
}"#;
storage
.write_byte_stream(&paths.lineage(), lineage_json.as_bytes().to_vec().into())
.await?;
let domain_lineage_io = DomainLineageIo::new(paths.lineage());
let package = InstalledPackage {
lineage: PackageLineageIo::new(domain_lineage_io, namespace.clone()),
paths,
remote,
storage,
namespace,
};
package
.set_remote(
"my-bucket".to_string(),
Some("my.host".parse()?),
WorkflowIntent::BucketDefault,
)
.await?;
let lineage = package.lineage().await?;
let remote_uri = lineage
.remote_uri
.as_ref()
.expect("remote_uri should be set");
assert_eq!(remote_uri.hash, "abc123", "hash should be preserved");
Ok(())
}
#[test(tokio::test)]
async fn test_set_remote_overwrites_unpushed_remote() -> Res {
let (home, _temp_dir1) = Home::from_temp_dir()?;
let (paths, _temp_dir2) = DomainPaths::from_temp_dir()?;
let storage = LocalStorage::new();
let remote = MockRemote::default();
let namespace: Namespace = ("test", "unpushed").into();
paths
.scaffold_for_installing(&storage, &home, &namespace)
.await?;
let lineage_json = r#"{
"packages": {
"test/unpushed": {
"commit": null,
"remote": {
"bucket": "old-bucket",
"namespace": "test/unpushed",
"hash": "",
"origin": "old.host"
},
"base_hash": "",
"latest_hash": "",
"paths": {}
}
},
"home": "/tmp/working_dir"
}"#;
storage
.write_byte_stream(&paths.lineage(), lineage_json.as_bytes().to_vec().into())
.await?;
let domain_lineage_io = DomainLineageIo::new(paths.lineage());
let package = InstalledPackage {
lineage: PackageLineageIo::new(domain_lineage_io, namespace.clone()),
paths,
remote,
storage,
namespace,
};
package
.set_remote(
"new-bucket".to_string(),
Some("new.host".parse()?),
WorkflowIntent::BucketDefault,
)
.await?;
let lineage = package.lineage().await?;
let remote_uri = lineage
.remote_uri
.as_ref()
.expect("remote_uri should be set");
assert_eq!(remote_uri.origin.as_ref().unwrap().to_string(), "new.host");
assert_eq!(remote_uri.bucket, "new-bucket");
assert_eq!(remote_uri.hash, "", "hash should remain empty");
Ok(())
}
#[test(tokio::test)]
async fn test_set_remote_recommits_existing_commit() -> Res {
let (home, _temp_dir1) = Home::from_temp_dir()?;
let (paths, _temp_dir2) = DomainPaths::from_temp_dir()?;
let storage = LocalStorage::new();
let remote = MockRemote::default();
let namespace: Namespace = ("test", "recommit").into();
paths
.scaffold_for_installing(&storage, &home, &namespace)
.await?;
let lineage_json = r#"{
"packages": {
"test/recommit": {
"commit": null,
"remote": null,
"base_hash": "",
"latest_hash": "",
"paths": {}
}
},
"home": "/tmp/working_dir"
}"#;
storage
.write_byte_stream(&paths.lineage(), lineage_json.as_bytes().to_vec().into())
.await?;
let package_home = home.join(namespace.to_string());
storage.create_dir_all(&package_home).await?;
storage
.write_byte_stream(
package_home.join("data.txt"),
ByteStream::from_static(b"hello world"),
)
.await?;
let domain_lineage_io = DomainLineageIo::new(paths.lineage());
let package = InstalledPackage {
lineage: PackageLineageIo::new(domain_lineage_io, namespace.clone()),
paths,
remote,
storage,
namespace: namespace.clone(),
};
let commit = package
.commit(
"Initial commit".to_string(),
UserMeta::Set(serde_json::json!({"key": "value"})),
None,
None,
)
.await?;
let hash_before = commit.hash.clone();
package
.set_remote(
"my-bucket".to_string(),
Some("example.com".parse()?),
WorkflowIntent::BucketDefault,
)
.await?;
let lineage = package.lineage().await?;
let remote_uri = lineage
.remote_uri
.as_ref()
.expect("remote_uri should be set");
assert_eq!(
remote_uri.origin.as_ref().unwrap().to_string(),
"example.com"
);
assert_eq!(remote_uri.bucket, "my-bucket");
let new_commit = lineage.commit.as_ref().expect("commit should exist");
assert_eq!(
new_commit.prev_hashes.first(),
Some(&hash_before),
"Old hash should be in prev_hashes after recommit"
);
let manifest_path = package
.paths
.installed_manifest(&namespace, &new_commit.hash);
let manifest = Manifest::from_path(&package.storage, &manifest_path).await?;
assert_eq!(
manifest.header.message,
Some("Initial commit".to_string()),
"Message should be preserved after recommit"
);
assert_eq!(
manifest.header.user_meta,
Some(serde_json::json!({"key": "value"})),
"User meta should be preserved after recommit"
);
Ok(())
}
#[test(tokio::test)]
async fn test_resolve_workflow_without_remote_is_none_for_every_intent() -> Res {
let (home, _temp_dir1) = Home::from_temp_dir()?;
let (paths, _temp_dir2) = DomainPaths::from_temp_dir()?;
let storage = LocalStorage::new();
let remote = MockRemote::default();
let namespace: Namespace = ("test", "noremote").into();
paths
.scaffold_for_installing(&storage, &home, &namespace)
.await?;
let lineage_json = r#"{
"packages": {
"test/noremote": {
"commit": null,
"remote": null,
"base_hash": "",
"latest_hash": "",
"paths": {}
}
},
"home": "/tmp/working_dir"
}"#;
storage
.write_byte_stream(&paths.lineage(), lineage_json.as_bytes().to_vec().into())
.await?;
let domain_lineage_io = DomainLineageIo::new(paths.lineage());
let package = InstalledPackage {
lineage: PackageLineageIo::new(domain_lineage_io, namespace.clone()),
paths,
remote,
storage,
namespace,
};
for intent in [
WorkflowIntent::BucketDefault,
WorkflowIntent::NoWorkflow,
WorkflowIntent::Named("foo".to_string()),
] {
assert!(
package.resolve_workflow(intent.clone()).await?.is_none(),
"no-remote short-circuit should return None for {intent:?}"
);
}
Ok(())
}
#[test(tokio::test)]
async fn test_set_remote_recommit_picks_up_bucket_default() -> Res {
let (home, _temp_dir1) = Home::from_temp_dir()?;
let (paths, _temp_dir2) = DomainPaths::from_temp_dir()?;
let storage = LocalStorage::new();
let remote = MockRemote::default();
let namespace: Namespace = ("test", "bucketdefault").into();
paths
.scaffold_for_installing(&storage, &home, &namespace)
.await?;
let config_uri: S3Uri = "s3://my-bucket/.quilt/workflows/config.yml".parse()?;
let config = r"
version: '1'
default_workflow: foo
workflows:
foo:
name: Foo
metadata_schema: bar
schemas:
bar:
url: s3://my-bucket/schemas/test.json
";
let schema_uri: S3Uri = "s3://my-bucket/schemas/test.json".parse()?;
remote
.put_object(None, &config_uri, config.as_bytes().to_vec())
.await?;
remote.put_object(None, &schema_uri, b"{}".to_vec()).await?;
let lineage_json = r#"{
"packages": {
"test/bucketdefault": {
"commit": null,
"remote": null,
"base_hash": "",
"latest_hash": "",
"paths": {}
}
},
"home": "/tmp/working_dir"
}"#;
storage
.write_byte_stream(&paths.lineage(), lineage_json.as_bytes().to_vec().into())
.await?;
let package_home = home.join(namespace.to_string());
storage.create_dir_all(&package_home).await?;
storage
.write_byte_stream(
package_home.join("data.txt"),
ByteStream::from_static(b"hello world"),
)
.await?;
let domain_lineage_io = DomainLineageIo::new(paths.lineage());
let package = InstalledPackage {
lineage: PackageLineageIo::new(domain_lineage_io, namespace.clone()),
paths,
remote,
storage,
namespace: namespace.clone(),
};
package
.commit(
"Initial commit".to_string(),
UserMeta::Set(serde_json::json!({"key": "value"})),
None,
None,
)
.await?;
let outcome = package
.set_remote(
"my-bucket".to_string(),
Some("example.com".parse()?),
WorkflowIntent::BucketDefault,
)
.await?;
assert!(
outcome.resolution_warning.is_none(),
"a clean bucket-default resolution must not produce a warning"
);
let lineage = package.lineage().await?;
let new_commit = lineage.commit.as_ref().expect("commit should exist");
let manifest_path = package
.paths
.installed_manifest(&namespace, &new_commit.hash);
let manifest = Manifest::from_path(&package.storage, &manifest_path).await?;
let workflow = manifest
.header
.workflow
.expect("recommit should stamp a workflow from the bucket default");
assert_eq!(
workflow.id.expect("workflow id should be set").id,
"foo",
"recommit should pick up the bucket's default_workflow"
);
Ok(())
}
#[test(tokio::test)]
async fn test_set_remote_fetches_config_and_schema_once() -> Res {
let (home, _temp_dir1) = Home::from_temp_dir()?;
let (paths, _temp_dir2) = DomainPaths::from_temp_dir()?;
let storage = LocalStorage::new();
let remote = MockRemote::default();
let namespace: Namespace = ("test", "fetchonce").into();
paths
.scaffold_for_installing(&storage, &home, &namespace)
.await?;
let config_uri_str = "s3://my-bucket/.quilt/workflows/config.yml";
let schema_uri_str = "s3://my-bucket/schemas/test.json";
let config_uri: S3Uri = config_uri_str.parse()?;
let config = r"
version: '1'
default_workflow: foo
workflows:
foo:
name: Foo
metadata_schema: bar
schemas:
bar:
url: s3://my-bucket/schemas/test.json
";
let schema_uri: S3Uri = schema_uri_str.parse()?;
remote
.put_object(None, &config_uri, config.as_bytes().to_vec())
.await?;
remote.put_object(None, &schema_uri, b"{}".to_vec()).await?;
let lineage_json = r#"{
"packages": {
"test/fetchonce": {
"commit": null,
"remote": null,
"base_hash": "",
"latest_hash": "",
"paths": {}
}
},
"home": "/tmp/working_dir"
}"#;
storage
.write_byte_stream(&paths.lineage(), lineage_json.as_bytes().to_vec().into())
.await?;
let package_home = home.join(namespace.to_string());
storage.create_dir_all(&package_home).await?;
storage
.write_byte_stream(
package_home.join("data.txt"),
ByteStream::from_static(b"hello world"),
)
.await?;
let domain_lineage_io = DomainLineageIo::new(paths.lineage());
let package = InstalledPackage {
lineage: PackageLineageIo::new(domain_lineage_io, namespace.clone()),
paths,
remote,
storage,
namespace: namespace.clone(),
};
package
.commit(
"Initial commit".to_string(),
UserMeta::Set(serde_json::json!({"key": "value"})),
None,
None,
)
.await?;
package
.set_remote(
"my-bucket".to_string(),
Some("example.com".parse()?),
WorkflowIntent::BucketDefault,
)
.await?;
assert_eq!(
package.remote.get_object_count(config_uri_str),
1,
"config.yml must be fetched exactly once across set_remote"
);
assert_eq!(
package.remote.get_object_count(schema_uri_str),
1,
"the schema document must be fetched exactly once across set_remote"
);
Ok(())
}
async fn recommit_manifest_for_intent(slug: &str, intent: WorkflowIntent) -> Res<ManifestHeader> {
let (home, _temp_dir1) = Home::from_temp_dir()?;
let (paths, _temp_dir2) = DomainPaths::from_temp_dir()?;
let storage = LocalStorage::new();
let remote = MockRemote::default();
let namespace: Namespace = ("test", slug).into();
paths
.scaffold_for_installing(&storage, &home, &namespace)
.await?;
let config_uri: S3Uri = "s3://my-bucket/.quilt/workflows/config.yml".parse()?;
let config = r"
version: '1'
is_workflow_required: false
workflows:
foo:
name: Foo
metadata_schema: bar
schemas:
bar:
url: s3://my-bucket/schemas/test.json
";
let schema_uri: S3Uri = "s3://my-bucket/schemas/test.json".parse()?;
remote
.put_object(None, &config_uri, config.as_bytes().to_vec())
.await?;
remote.put_object(None, &schema_uri, b"{}".to_vec()).await?;
let lineage_json = format!(
r#"{{
"packages": {{
"test/{slug}": {{
"commit": null,
"remote": null,
"base_hash": "",
"latest_hash": "",
"paths": {{}}
}}
}},
"home": "/tmp/working_dir"
}}"#
);
storage
.write_byte_stream(&paths.lineage(), lineage_json.as_bytes().to_vec().into())
.await?;
let package_home = home.join(namespace.to_string());
storage.create_dir_all(&package_home).await?;
storage
.write_byte_stream(
package_home.join("data.txt"),
ByteStream::from_static(b"hello world"),
)
.await?;
let domain_lineage_io = DomainLineageIo::new(paths.lineage());
let package = InstalledPackage {
lineage: PackageLineageIo::new(domain_lineage_io, namespace.clone()),
paths,
remote,
storage,
namespace: namespace.clone(),
};
package
.commit(
"Initial commit".to_string(),
UserMeta::Set(serde_json::json!({"key": "value"})),
None,
None,
)
.await?;
package
.set_remote(
"my-bucket".to_string(),
Some("example.com".parse()?),
intent,
)
.await?;
let lineage = package.lineage().await?;
let new_commit = lineage.commit.as_ref().expect("commit should exist");
let manifest_path = package
.paths
.installed_manifest(&namespace, &new_commit.hash);
let manifest = Manifest::from_path(&package.storage, &manifest_path).await?;
Ok(manifest.header)
}
const FOO_CONFIG: &str = r"
version: '1'
workflows:
foo:
name: Foo
metadata_schema: bar
schemas:
bar:
url: s3://my-bucket/schemas/test.json
";
async fn package_with_config(
slug: &str,
config: &str,
schema: &[u8],
) -> Res<(
InstalledPackage<LocalStorage, MockRemote>,
tempfile::TempDir,
tempfile::TempDir,
)> {
let (home, temp_dir1) = Home::from_temp_dir()?;
let (paths, temp_dir2) = DomainPaths::from_temp_dir()?;
let storage = LocalStorage::new();
let remote = MockRemote::default();
let namespace: Namespace = ("test", slug).into();
paths
.scaffold_for_installing(&storage, &home, &namespace)
.await?;
let config_uri: S3Uri = "s3://my-bucket/.quilt/workflows/config.yml".parse()?;
let schema_uri: S3Uri = "s3://my-bucket/schemas/test.json".parse()?;
remote
.put_object(None, &config_uri, config.as_bytes().to_vec())
.await?;
remote
.put_object(None, &schema_uri, schema.to_vec())
.await?;
let lineage_json = format!(
r#"{{
"packages": {{
"test/{slug}": {{
"commit": null,
"remote": null,
"base_hash": "",
"latest_hash": "",
"paths": {{}}
}}
}},
"home": "/tmp/working_dir"
}}"#
);
storage
.write_byte_stream(&paths.lineage(), lineage_json.as_bytes().to_vec().into())
.await?;
let package_home = home.join(namespace.to_string());
storage.create_dir_all(&package_home).await?;
storage
.write_byte_stream(
package_home.join("data.txt"),
ByteStream::from_static(b"hello world"),
)
.await?;
let domain_lineage_io = DomainLineageIo::new(paths.lineage());
let package = InstalledPackage {
lineage: PackageLineageIo::new(domain_lineage_io, namespace.clone()),
paths,
remote,
storage,
namespace,
};
package
.commit(
"Initial commit".to_string(),
UserMeta::Set(serde_json::json!({"key": "value"})),
None,
None,
)
.await?;
Ok((package, temp_dir1, temp_dir2))
}
#[test(tokio::test)]
async fn test_set_remote_propagates_named_workflow_error() -> Res {
let (package, _t1, _t2) = package_with_config("named-error", FOO_CONFIG, b"{}").await?;
let result = package
.set_remote(
"my-bucket".to_string(),
Some("example.com".parse()?),
WorkflowIntent::Named("nope".to_string()),
)
.await;
assert!(
result.is_err(),
"an explicit Named intent with an unknown id must surface the recommit error"
);
assert!(
result.unwrap_err().to_string().contains("Workflow nope"),
"error should name the unresolved workflow"
);
Ok(())
}
#[test(tokio::test)]
async fn test_set_remote_swallows_bucket_default_recommit_error() -> Res {
let config = r"
version: '1'
default_workflow: ghost
workflows:
foo:
name: Foo
metadata_schema: bar
schemas:
bar:
url: s3://my-bucket/schemas/test.json
";
let (package, _t1, _t2) = package_with_config("bucketdefault-ok", config, b"{}").await?;
let outcome = package
.set_remote(
"my-bucket".to_string(),
Some("example.com".parse()?),
WorkflowIntent::BucketDefault,
)
.await?;
let warning = outcome
.resolution_warning
.expect("a failed bucket-default resolution must surface a warning");
assert!(
warning.contains("ghost"),
"the warning must carry the underlying reason, got: {warning}"
);
assert!(
!warning.contains("Remote catalog error"),
"the warning must be the unwrapped inner reason, not the error chain, got: {warning}"
);
let lineage = package.lineage().await?;
assert!(
lineage.remote_uri.is_some(),
"remote should be persisted on the BucketDefault path"
);
let commit = lineage.commit.as_ref().expect("commit should still exist");
let manifest_path = package
.paths
.installed_manifest(&package.namespace, &commit.hash);
let manifest = Manifest::from_path(&package.storage, &manifest_path).await?;
assert!(
manifest.header.workflow.is_none(),
"no workflow may be stamped when the bucket default fails to resolve"
);
Ok(())
}
async fn assert_nothing_persisted(
package: &InstalledPackage<LocalStorage, MockRemote>,
hash_before: &str,
manifests_before: usize,
) -> Res {
let lineage = package.lineage().await?;
assert!(
lineage.remote_uri.is_none(),
"a rejected set_remote must not persist the remote"
);
let commit = lineage.commit.as_ref().expect("commit should still exist");
assert_eq!(
commit.hash, hash_before,
"a rejected set_remote must not change the commit"
);
assert!(
commit.prev_hashes.is_empty(),
"a rejected set_remote must not record a recommit in prev_hashes"
);
let manifests_dir = package.paths.installed_manifests_dir(&package.namespace);
let manifest_path = package
.paths
.installed_manifest(&package.namespace, hash_before);
let manifest = Manifest::from_path(&package.storage, &manifest_path).await?;
assert!(
manifest.header.workflow.is_none(),
"the previous manifest's header must be unchanged"
);
let manifests_after = std::fs::read_dir(&manifests_dir)?.count();
assert_eq!(
manifests_after, manifests_before,
"a rejected set_remote must not write a new manifest"
);
Ok(())
}
#[test(tokio::test)]
async fn test_set_remote_no_workflow_against_required_bucket_is_rejected() -> Res {
let (package, _t1, _t2) = package_with_config("noworkflow-required", FOO_CONFIG, b"{}").await?;
let lineage = package.lineage().await?;
let hash_before = lineage.commit.as_ref().expect("committed").hash.clone();
let manifests_dir = package.paths.installed_manifests_dir(&package.namespace);
let manifests_before = std::fs::read_dir(&manifests_dir)?.count();
let err = package
.set_remote(
"my-bucket".to_string(),
Some("example.com".parse()?),
WorkflowIntent::NoWorkflow,
)
.await
.unwrap_err();
assert!(
matches!(
&err,
Error::WorkflowValidation(WorkflowValidationError::Rejected(violations))
if violations.contains(&RuleViolation::WorkflowRequired)
),
"expected a WorkflowRequired rejection, got: {err:?}"
);
assert_nothing_persisted(&package, &hash_before, manifests_before).await?;
Ok(())
}
#[test(tokio::test)]
async fn test_set_remote_bucket_default_validation_error_propagates() -> Res {
let config = r"
version: '1'
default_workflow: foo
workflows:
foo:
name: Foo
metadata_schema: bar
schemas:
bar:
url: s3://my-bucket/schemas/test.json
";
let schema = br#"{"type": "object", "required": ["owner"]}"#;
let (package, _t1, _t2) = package_with_config("bucketdefault-invalid", config, schema).await?;
let lineage = package.lineage().await?;
let hash_before = lineage.commit.as_ref().expect("committed").hash.clone();
let manifests_dir = package.paths.installed_manifests_dir(&package.namespace);
let manifests_before = std::fs::read_dir(&manifests_dir)?.count();
let err = package
.set_remote(
"my-bucket".to_string(),
Some("example.com".parse()?),
WorkflowIntent::BucketDefault,
)
.await
.unwrap_err();
assert!(
matches!(
&err,
Error::WorkflowValidation(WorkflowValidationError::Rejected(violations))
if matches!(&violations[..], [RuleViolation::MetadataInvalid(_)])
),
"expected a MetadataInvalid rejection, got: {err:?}"
);
assert_nothing_persisted(&package, &hash_before, manifests_before).await?;
Ok(())
}
#[test(tokio::test)]
async fn test_set_remote_stamps_named_workflow() -> Res {
let header =
recommit_manifest_for_intent("named", WorkflowIntent::Named("foo".to_string())).await?;
let workflow = header
.workflow
.expect("recommit should stamp the named workflow");
assert_eq!(
workflow.id.expect("workflow id should be set").id,
"foo",
"recommit should stamp the caller's chosen workflow, not the bucket default"
);
Ok(())
}
#[test(tokio::test)]
async fn test_set_remote_stamps_no_workflow() -> Res {
let header = recommit_manifest_for_intent("noworkflow", WorkflowIntent::NoWorkflow).await?;
let workflow = header
.workflow
.expect("recommit should stamp an id-less workflow when a config is present");
assert!(
workflow.id.is_none(),
"NoWorkflow must not resolve any workflow id"
);
Ok(())
}