use std::path::Path;
use std::{fs, io};
use crate::bootstrap_config::policy_store_config::{PolicyStoreConfig, PolicyStoreSource};
use crate::common::policy_store::errors::{PolicyStoreError, ValidationError};
use crate::common::policy_store::legacy_store::LegacyAgamaPolicyStore;
use crate::common::policy_store::manager::PolicyStoreManager;
use crate::common::policy_store::validator::MetadataValidator;
use crate::common::policy_store::{ConversionError, PolicyStore, PolicyStoreWithID};
use crate::http::cache_headers::CacheHeadersState;
use crate::http::{HttpClient, HttpClientError};
pub(super) const ZIP_MAGIC: [u8; 4] = [0x50, 0x4B, 0x03, 0x04];
#[derive(Debug, thiserror::Error)]
pub enum PolicyStoreLoadError {
#[error("failed to parse the policy store from policy_store json: {0}")]
ParseJson(#[from] serde_json::Error),
#[error("failed to parse the policy store from policy_store yaml: {0}")]
ParseYaml(#[from] serde_yaml_ng::Error),
#[error("failed to fetch the policy store from the lock server")]
FetchFromLockServer(#[from] HttpClientError),
#[error("Policy Store does not contain correct structure: {0}")]
InvalidStore(String),
#[error("Failed to load policy store from {0}: {1}")]
ParseFile(Box<Path>, io::Error),
#[error("Failed to convert loaded policy store: {0}")]
Conversion(#[from] ConversionError),
#[error("Failed to load policy store from archive: {0}")]
Archive(String),
#[error("Failed to load policy store from directory: {0}")]
Directory(String),
#[error("Policy store validation error: {0}")]
Validation(#[from] ValidationError),
}
fn extract_first_policy_store(
agama_policy_store: &LegacyAgamaPolicyStore,
strict_schema_validation: bool,
) -> Result<PolicyStoreWithID, PolicyStoreLoadError> {
if agama_policy_store.policy_stores.len() != 1 {
return Err(PolicyStoreLoadError::InvalidStore(format!(
"expected exactly one 'policy_stores' entry, but found {:?}",
agama_policy_store.policy_stores.len()
)));
}
let policy_store_option = agama_policy_store
.policy_stores
.iter()
.take(1)
.map(|(k, v)| {
let store: PolicyStore = v.to_owned().into();
let metadata = crate::common::policy_store::metadata::PolicyStoreMetadata {
cedar_version: agama_policy_store
.cedar_version
.strip_prefix('v')
.unwrap_or(&agama_policy_store.cedar_version)
.to_string(),
policy_store: crate::common::policy_store::metadata::PolicyStoreInfo {
id: String::new(),
name: k.clone(),
version: v.version.clone().unwrap_or_default(),
description: None,
created_date: None,
updated_date: None,
},
};
PolicyStoreWithID {
id: k.to_owned(),
store,
metadata: Some(metadata),
}
})
.next();
match policy_store_option {
Some(policy_store) => {
if strict_schema_validation && policy_store.schema.is_none() {
return Err(PolicyStoreLoadError::InvalidStore(
"missing required schema in policy store".to_string(),
));
}
Ok(policy_store)
},
None => Err(PolicyStoreLoadError::InvalidStore(
"error retrieving first policy_stores element".into(),
)),
}
}
pub(crate) struct LoadedPolicyStore {
pub store: PolicyStoreWithID,
pub body_hash: Option<u64>,
pub validators: CacheHeadersState,
}
pub(crate) async fn load_policy_store(
config: &PolicyStoreConfig,
http_client: &HttpClient,
strict_schema_validation: bool,
) -> Result<LoadedPolicyStore, PolicyStoreLoadError> {
let loaded = match &config.source {
PolicyStoreSource::Json(policy_json) => {
let agama_policy_store = serde_json::from_str::<LegacyAgamaPolicyStore>(policy_json)
.map_err(PolicyStoreLoadError::ParseJson)?;
process_legacy_agama_store(&agama_policy_store, strict_schema_validation)?
},
PolicyStoreSource::Yaml(policy_yaml) => {
let agama_policy_store = serde_yaml_ng::from_str::<LegacyAgamaPolicyStore>(policy_yaml)
.map_err(PolicyStoreLoadError::ParseYaml)?;
process_legacy_agama_store(&agama_policy_store, strict_schema_validation)?
},
PolicyStoreSource::LockServer(policy_store_uri) => {
load_policy_store_from_lock_master(
policy_store_uri,
http_client,
strict_schema_validation,
)
.await?
},
PolicyStoreSource::FileJson(path) => {
let policy_json = fs::read_to_string(path)
.map_err(|e| PolicyStoreLoadError::ParseFile(path.clone().into(), e))?;
let agama_policy_store = serde_json::from_str::<LegacyAgamaPolicyStore>(&policy_json)
.map_err(PolicyStoreLoadError::ParseJson)?;
process_legacy_agama_store(&agama_policy_store, strict_schema_validation)?
},
PolicyStoreSource::FileYaml(path) => {
let policy_yaml = fs::read_to_string(path)
.map_err(|e| PolicyStoreLoadError::ParseFile(path.clone().into(), e))?;
let agama_policy_store =
serde_yaml_ng::from_str::<LegacyAgamaPolicyStore>(&policy_yaml)
.map_err(PolicyStoreLoadError::ParseYaml)?;
process_legacy_agama_store(&agama_policy_store, strict_schema_validation)?
},
#[cfg(not(target_arch = "wasm32"))]
PolicyStoreSource::CjarFile(path) => LoadedPolicyStore {
store: load_policy_store_from_cjar_file(path, strict_schema_validation).await?,
body_hash: None,
validators: CacheHeadersState::default(),
},
#[cfg(target_arch = "wasm32")]
PolicyStoreSource::CjarFile(path) => LoadedPolicyStore {
store: load_policy_store_from_cjar_file(path, strict_schema_validation)?,
body_hash: None,
validators: CacheHeadersState::default(),
},
PolicyStoreSource::CjarUrl(url) => {
load_policy_store_from_cjar_url(url, http_client, strict_schema_validation).await?
},
#[cfg(not(target_arch = "wasm32"))]
PolicyStoreSource::Directory(path) => LoadedPolicyStore {
store: load_policy_store_from_directory(path, strict_schema_validation).await?,
body_hash: None,
validators: CacheHeadersState::default(),
},
#[cfg(target_arch = "wasm32")]
PolicyStoreSource::Directory(path) => LoadedPolicyStore {
store: load_policy_store_from_directory(path, strict_schema_validation)?,
body_hash: None,
validators: CacheHeadersState::default(),
},
PolicyStoreSource::ArchiveBytes(bytes) => LoadedPolicyStore {
store: load_policy_store_from_archive_bytes(bytes, strict_schema_validation)?,
body_hash: None,
validators: CacheHeadersState::default(),
},
PolicyStoreSource::Uri(uri) => {
load_policy_store_from_uri(uri, http_client, strict_schema_validation).await?
},
};
Ok(loaded)
}
fn process_legacy_agama_store(
agama_policy_store: &LegacyAgamaPolicyStore,
strict_schema_validation: bool,
) -> Result<LoadedPolicyStore, PolicyStoreLoadError> {
MetadataValidator::validate_legacy_store(agama_policy_store)
.map_err(PolicyStoreLoadError::Validation)?;
Ok(LoadedPolicyStore {
store: extract_first_policy_store(agama_policy_store, strict_schema_validation)?,
body_hash: None,
validators: CacheHeadersState::default(),
})
}
async fn load_policy_store_from_uri(
uri: &str,
http_client: &HttpClient,
strict_schema_validation: bool,
) -> Result<LoadedPolicyStore, PolicyStoreLoadError> {
let response = http_client.get_with_retry(uri).await?;
let validators = CacheHeadersState::from_headers(response.headers(), chrono::Utc::now());
let bytes = http_client.read_response_capped(response).await?;
let body_hash = crate::init::policy_store_refresh::body_hash(&bytes);
if bytes.starts_with(&ZIP_MAGIC) {
return Ok(LoadedPolicyStore {
store: parse_cjar_bytes(&bytes, strict_schema_validation).await?,
body_hash: Some(body_hash),
validators,
});
}
let store = parse_lock_master_bytes(&bytes, strict_schema_validation)?;
Ok(LoadedPolicyStore {
store,
body_hash: Some(body_hash),
validators,
})
}
async fn load_policy_store_from_lock_master(
uri: &str,
http_client: &HttpClient,
strict_schema_validation: bool,
) -> Result<LoadedPolicyStore, PolicyStoreLoadError> {
let response = http_client.get_with_retry(uri).await?;
let validators = CacheHeadersState::from_headers(response.headers(), chrono::Utc::now());
let bytes = http_client.read_response_capped(response).await?;
let store = parse_lock_master_bytes(&bytes, strict_schema_validation)?;
Ok(LoadedPolicyStore {
store,
body_hash: Some(crate::init::policy_store_refresh::body_hash(&bytes)),
validators,
})
}
pub(crate) fn parse_lock_master_bytes(
bytes: &[u8],
strict_schema_validation: bool,
) -> Result<PolicyStoreWithID, PolicyStoreLoadError> {
let agama_policy_store: LegacyAgamaPolicyStore = serde_json::from_slice(bytes)?;
let loaded = process_legacy_agama_store(&agama_policy_store, strict_schema_validation)?;
Ok(loaded.store)
}
pub(crate) async fn parse_cjar_bytes(
bytes: &[u8],
strict_schema_validation: bool,
) -> Result<PolicyStoreWithID, PolicyStoreLoadError> {
use crate::common::policy_store::loader;
let loaded = loader::load_policy_store_archive_bytes(bytes, strict_schema_validation)
.map_err(|e| PolicyStoreLoadError::Archive(format!("Failed to load from archive: {e}")))?;
let store_id = loaded.metadata.policy_store.id.clone();
let store_metadata = loaded.metadata.clone();
#[cfg(not(target_arch = "wasm32"))]
let legacy_store = tokio::task::spawn_blocking(move || {
convert_archive_to_legacy(loaded, strict_schema_validation)
})
.await
.map_err(|e| PolicyStoreLoadError::Archive(format!("Conversion task panicked: {e}")))??;
#[cfg(target_arch = "wasm32")]
let legacy_store = convert_archive_to_legacy(loaded, strict_schema_validation)?;
Ok(PolicyStoreWithID {
id: store_id,
store: legacy_store,
metadata: Some(store_metadata),
})
}
fn convert_archive_to_legacy(
loaded: crate::common::policy_store::loader::LoadedPolicyStore,
strict_schema_validation: bool,
) -> Result<crate::common::policy_store::PolicyStore, PolicyStoreLoadError> {
PolicyStoreManager::convert_to_legacy(loaded, strict_schema_validation).map_err(Into::into)
}
#[cfg(not(target_arch = "wasm32"))]
async fn load_policy_store_from_cjar_file(
path: &Path,
strict_schema_validation: bool,
) -> Result<PolicyStoreWithID, PolicyStoreLoadError> {
use crate::common::policy_store::loader;
let loaded = loader::load_policy_store_archive(path, strict_schema_validation)
.await
.map_err(|e| map_policy_store_err(e, true))?;
let store_id = loaded.metadata.policy_store.id.clone();
let store_metadata = loaded.metadata.clone();
let legacy_store = tokio::task::spawn_blocking(move || {
convert_archive_to_legacy(loaded, strict_schema_validation)
})
.await
.map_err(|e| PolicyStoreLoadError::Archive(format!("Conversion task panicked: {e}")))??;
Ok(PolicyStoreWithID {
id: store_id,
store: legacy_store,
metadata: Some(store_metadata),
})
}
#[cfg(target_arch = "wasm32")]
fn load_policy_store_from_cjar_file(
path: &Path,
strict_schema_validation: bool,
) -> Result<PolicyStoreWithID, PolicyStoreLoadError> {
use crate::common::policy_store::loader;
match loader::load_policy_store_archive(path, strict_schema_validation) {
Err(e) => Err(PolicyStoreLoadError::Archive(format!(
"Loading from file path is not supported in WASM. Use CjarUrl instead. Original error: {e}",
))),
Ok(_) => unreachable!("WASM stub should always return an error"),
}
}
async fn load_policy_store_from_cjar_url(
url: &str,
http_client: &HttpClient,
strict_schema_validation: bool,
) -> Result<LoadedPolicyStore, PolicyStoreLoadError> {
use crate::common::policy_store::loader;
let response = http_client
.get_with_retry(url)
.await
.map_err(|e| PolicyStoreLoadError::Archive(format!("Failed to fetch archive: {e}")))?;
let validators = CacheHeadersState::from_headers(response.headers(), chrono::Utc::now());
let bytes = http_client
.read_response_capped(response)
.await
.map_err(|e| PolicyStoreLoadError::Archive(format!("Failed to read archive body: {e}")))?;
let body_hash = crate::init::policy_store_refresh::body_hash(&bytes);
let loaded = loader::load_policy_store_archive_bytes(&bytes, strict_schema_validation)
.map_err(|e| map_policy_store_err(e, true))?;
let store_id = loaded.metadata.policy_store.id.clone();
let store_metadata = loaded.metadata.clone();
#[cfg(not(target_arch = "wasm32"))]
let legacy_store = tokio::task::spawn_blocking(move || {
convert_archive_to_legacy(loaded, strict_schema_validation)
})
.await
.map_err(|e| PolicyStoreLoadError::Archive(format!("Conversion task panicked: {e}")))??;
#[cfg(target_arch = "wasm32")]
let legacy_store = convert_archive_to_legacy(loaded, strict_schema_validation)?;
Ok(LoadedPolicyStore {
store: PolicyStoreWithID {
id: store_id,
store: legacy_store,
metadata: Some(store_metadata),
},
body_hash: Some(body_hash),
validators,
})
}
#[cfg(not(target_arch = "wasm32"))]
async fn load_policy_store_from_directory(
path: &Path,
strict_schema_validation: bool,
) -> Result<PolicyStoreWithID, PolicyStoreLoadError> {
use crate::common::policy_store::loader;
let loaded = loader::load_policy_store_directory(path, strict_schema_validation)
.await
.map_err(|e| map_policy_store_err(e, false))?;
let store_id = loaded.metadata.policy_store.id.clone();
let store_metadata = loaded.metadata.clone();
let legacy_store = tokio::task::spawn_blocking(move || {
convert_archive_to_legacy(loaded, strict_schema_validation)
})
.await
.map_err(|e| PolicyStoreLoadError::Directory(format!("Conversion task panicked: {e}")))??;
Ok(PolicyStoreWithID {
id: store_id,
store: legacy_store,
metadata: Some(store_metadata),
})
}
#[cfg(target_arch = "wasm32")]
fn load_policy_store_from_directory(
path: &Path,
strict_schema_validation: bool,
) -> Result<PolicyStoreWithID, PolicyStoreLoadError> {
use crate::common::policy_store::loader;
match loader::load_policy_store_directory(path, strict_schema_validation) {
Err(e) => Err(PolicyStoreLoadError::Directory(format!(
"Loading from directory is not supported in WASM. Original error: {e}",
))),
Ok(_) => unreachable!("WASM stub should always return an error"),
}
}
fn load_policy_store_from_archive_bytes(
bytes: &[u8],
strict_schema_validation: bool,
) -> Result<PolicyStoreWithID, PolicyStoreLoadError> {
use crate::common::policy_store::loader;
let loaded =
loader::load_policy_store_archive_bytes(bytes, strict_schema_validation)
.map_err(|e| map_policy_store_err(e, true))?;
let store_id = loaded.metadata.policy_store.id.clone();
let store_metadata = loaded.metadata.clone();
let legacy_store = convert_archive_to_legacy(loaded, strict_schema_validation)?;
Ok(PolicyStoreWithID {
id: store_id,
store: legacy_store,
metadata: Some(store_metadata),
})
}
fn map_policy_store_err(e: PolicyStoreError, is_archive: bool) -> PolicyStoreLoadError {
match e {
PolicyStoreError::Validation(ve) => PolicyStoreLoadError::Validation(ve),
PolicyStoreError::CedarParsing { file, detail } => {
PolicyStoreLoadError::InvalidStore(format!("Cedar parse error in {file}: {detail}"))
},
PolicyStoreError::CedarSchemaError { file, err } => {
PolicyStoreLoadError::InvalidStore(format!("Cedar schema error in {file}: {err}"))
},
_ => {
if is_archive {
PolicyStoreLoadError::Archive(e.to_string())
} else {
PolicyStoreLoadError::Directory(e.to_string())
}
},
}
}
#[cfg(test)]
mod test {
use std::{path::Path, sync::LazyLock, time::Duration};
use base64::Engine;
use mockito::Server;
use serde_json::json;
use super::{extract_first_policy_store, load_policy_store};
use crate::common::policy_store::legacy_store::LegacyAgamaPolicyStore;
use crate::{
PolicyStoreConfig, PolicyStoreSource,
common::policy_store::test_utils::fixtures,
http::{HttpClient, HttpClientConfig},
};
static HTTP_CLIENT: LazyLock<HttpClient> = LazyLock::new(|| {
HttpClient::new(HttpClientConfig {
max_retries: 0,
retry_delay: Duration::from_millis(3),
request_timeout: Duration::from_millis(500),
max_response_size_bytes: None,
})
.expect("http client should be constructed")
});
fn make_full_legacy_json() -> serde_json::Value {
let schema = base64::prelude::BASE64_STANDARD.encode(
r#"{
"Jans": {
"entityTypes": {},
"actions": {}
}
}"#,
);
json!({
"cedar_version": "v4.0.0",
"policy_stores": {
"test": {
"name": "test",
"schema": schema,
"policies": {}
}
}
})
}
fn make_no_schema_legacy_json() -> serde_json::Value {
json!({
"cedar_version": "v4.0.0",
"policy_stores": {
"test": {
"name": "test",
"policies": {}
}
}
})
}
fn make_null_schema_legacy_json() -> serde_json::Value {
json!({
"cedar_version": "v4.0.0",
"policy_stores": {
"test": {
"name": "test",
"schema": null,
"policies": {}
}
}
})
}
#[test]
fn test_extract_first_policy_store_with_schema_strict_true() {
let agama: LegacyAgamaPolicyStore = serde_json::from_value(make_full_legacy_json())
.expect("valid legacy store with schema");
let result = extract_first_policy_store(&agama, true);
result.expect("should succeed with schema and strict=true");
}
#[test]
fn test_extract_first_policy_store_with_schema_strict_false() {
let agama: LegacyAgamaPolicyStore = serde_json::from_value(make_full_legacy_json())
.expect("valid legacy store with schema");
let result = extract_first_policy_store(&agama, false);
result.expect("should succeed with schema and strict=false");
}
#[test]
fn test_extract_first_policy_store_missing_schema_strict_true() {
let agama: LegacyAgamaPolicyStore = serde_json::from_value(make_no_schema_legacy_json())
.expect("valid legacy store without schema");
let result = extract_first_policy_store(&agama, true);
let err = result.expect_err("should error when schema missing and strict=true");
assert!(
err.to_string().contains("missing required schema"),
"error should mention missing schema, got: {err}"
);
}
#[test]
fn test_extract_first_policy_store_missing_schema_strict_false() {
let agama: LegacyAgamaPolicyStore = serde_json::from_value(make_no_schema_legacy_json())
.expect("valid legacy store without schema");
let result = extract_first_policy_store(&agama, false);
result.expect("should succeed when schema missing and strict=false");
}
#[test]
fn test_extract_first_policy_store_null_schema_strict_true() {
let agama: LegacyAgamaPolicyStore = serde_json::from_value(make_null_schema_legacy_json())
.expect("valid legacy store with null schema");
let result = extract_first_policy_store(&agama, true);
let err = result.expect_err("should error when schema is null and strict=true");
assert!(
err.to_string().contains("missing required schema"),
"error should mention missing schema, got: {err}"
);
}
#[test]
fn test_extract_first_policy_store_null_schema_strict_false() {
let agama: LegacyAgamaPolicyStore = serde_json::from_value(make_null_schema_legacy_json())
.expect("valid legacy store with null schema");
let result = extract_first_policy_store(&agama, false);
result.expect("should succeed when schema is null and strict=false");
}
#[tokio::test]
async fn can_load_from_json_file() {
load_policy_store(
&PolicyStoreConfig {
source: crate::PolicyStoreSource::FileJson(
Path::new("../test_files/policy-store_generated.json").into(),
),
..Default::default()
},
&HTTP_CLIENT,
true,
)
.await
.expect("Should load policy store from JSON file");
}
#[tokio::test]
async fn can_load_from_yaml_file() {
load_policy_store(
&PolicyStoreConfig {
source: crate::PolicyStoreSource::FileYaml(
Path::new("../test_files/policy-store_ok.yaml").into(),
),
..Default::default()
},
&HTTP_CLIENT,
true,
)
.await
.expect("Should load policy store from YAML file");
}
#[tokio::test]
async fn can_load_from_lock_master() {
let mut mock_server = Server::new_async().await;
let policy_store_json =
include_str!("../../../test_files/policy-store_lock_master_ok.json").to_string();
let mock_endpoint = mock_server
.mock("GET", "/policy-store")
.with_status(200)
.with_header("content-type", "application/json")
.with_body(policy_store_json)
.expect(1)
.create();
let uri = format!("{}/policy-store", mock_server.url()).to_string();
load_policy_store(
&PolicyStoreConfig {
source: crate::PolicyStoreSource::LockServer(uri),
..Default::default()
},
&HTTP_CLIENT,
true,
)
.await
.expect("Should load policy store from Lock Master file");
mock_endpoint.assert();
}
#[tokio::test]
async fn can_load_from_uri_with_json_content_type() {
let mut mock_server = Server::new_async().await;
let policy_store_json =
include_str!("../../../test_files/policy-store_lock_master_ok.json").to_string();
let mock_endpoint = mock_server
.mock("GET", "/policy-store")
.with_status(200)
.with_header("content-type", "application/json")
.with_body(policy_store_json)
.expect(1)
.create();
let uri = format!("{}/policy-store", mock_server.url()).to_string();
load_policy_store(
&PolicyStoreConfig {
source: PolicyStoreSource::Uri(uri),
refresh_interval_secs: 0,
},
&HTTP_CLIENT,
false,
)
.await
.expect("Should load policy store from URI with JSON content-type");
mock_endpoint.assert();
}
#[tokio::test]
async fn can_load_from_uri_missing_content_type_uses_magic_bytes() {
let mut mock_server = Server::new_async().await;
let archive_bytes = fixtures::minimal_valid()
.build_archive()
.expect("Should build test archive");
let mock_endpoint = mock_server
.mock("GET", "/policy-store")
.with_status(200)
.with_body(archive_bytes)
.expect(1)
.create();
let uri = format!("{}/policy-store", mock_server.url()).to_string();
load_policy_store(
&PolicyStoreConfig {
source: PolicyStoreSource::Uri(uri),
refresh_interval_secs: 0,
},
&HTTP_CLIENT,
false,
)
.await
.expect("Should load policy store from URI with magic bytes and no content-type");
mock_endpoint.assert();
}
#[tokio::test]
async fn can_load_from_uri_with_octet_stream_content_type() {
let mut mock_server = Server::new_async().await;
let archive_bytes = fixtures::minimal_valid()
.build_archive()
.expect("Should build test archive");
let mock_endpoint = mock_server
.mock("GET", "/policy-store")
.with_status(200)
.with_header("content-type", "application/octet-stream")
.with_body(archive_bytes)
.expect(1)
.create();
let uri = format!("{}/policy-store", mock_server.url()).to_string();
load_policy_store(
&PolicyStoreConfig {
source: PolicyStoreSource::Uri(uri),
refresh_interval_secs: 0,
},
&HTTP_CLIENT,
false,
)
.await
.expect("Should load policy store from URI with octet-stream content-type");
mock_endpoint.assert();
}
}