use dataflow_rs::engine::error::DataflowError;
use dataflow_rs::engine::secrets::SECRET_OPERATOR;
use dataflow_rs::engine::task_context::TaskContext;
use serde_json::Value;
pub fn secret_name(value: &Value) -> Option<&str> {
let object = value.as_object()?;
if object.len() != 1 {
return None;
}
match object.get(SECRET_OPERATOR)? {
Value::String(name) => Some(name.as_str()),
Value::Array(items) if items.len() == 1 => items[0].as_str(),
_ => None,
}
}
async fn resolve_key_material(
value: &Value,
field: &str,
ctx: &TaskContext<'_>,
) -> Result<String, String> {
if let Some(name) = secret_name(value) {
let Some(secret) = ctx.secret(name) else {
return Err(format!(
"'{field}' names secret '{name}', which this instance does not declare \
— add it to the [secrets] section of the config file"
));
};
return secret
.as_str()
.map(str::to_string)
.ok_or_else(|| format!("'{field}': secret '{name}' is not a string"));
}
match value.as_str() {
Some(text) => crate::connector::secrets::resolve_secret_string(text, field).await,
None => Err(format!(
"'{field}' must be a string (a literal or an env:// reference) or \
{{\"secret\": \"name\"}}"
)),
}
}
pub async fn key_material_field(
input: &super::templated_input::TemplatedInput,
field: &str,
handler: &str,
ctx: &TaskContext<'_>,
) -> Result<String, DataflowError> {
let Some(authored) = input.get(field) else {
return Err(DataflowError::Validation(format!(
"{handler} requires '{field}' ({{\"secret\": \"name\"}}, a secret reference \
like env://NAME, an expression, or a literal)"
)));
};
let label = format!("{handler}.{field}");
if authored.is_string() {
return key_material(authored, &label, ctx).await;
}
let Some(value) = input.value_of(field, handler, ctx) else {
return Err(DataflowError::Validation(format!(
"{handler} requires '{field}'"
)));
};
match value? {
Value::String(material) => Ok(material),
other => Err(DataflowError::Validation(format!(
"'{label}' must resolve to a string, got {other}"
))),
}
}
pub async fn key_material(
value: &Value,
field: &str,
ctx: &TaskContext<'_>,
) -> Result<String, DataflowError> {
resolve_key_material(value, field, ctx)
.await
.map_err(DataflowError::Validation)
}
#[cfg(test)]
mod tests {
use super::*;
use serde_json::json;
#[test]
fn only_a_single_key_secret_object_is_a_reference() {
assert_eq!(secret_name(&json!({"secret": "k"})), Some("k"));
assert_eq!(secret_name(&json!({"secret": "k", "other": 1})), None);
assert_eq!(secret_name(&json!({"var": "data.k"})), None);
assert_eq!(secret_name(&json!("k")), None);
assert_eq!(secret_name(&json!({"secret": 7})), None);
}
#[test]
fn the_one_element_array_spelling_is_the_same_reference() {
assert_eq!(secret_name(&json!({"secret": ["k"]})), Some("k"));
assert_eq!(secret_name(&json!({"secret": ["k", "j"]})), None);
assert_eq!(secret_name(&json!({"secret": []})), None);
assert_eq!(secret_name(&json!({"secret": [7]})), None);
}
}