use crate::error::{CliError, CliResult};
use schemars::JsonSchema;
use serde::{Deserialize, Serialize};
#[derive(Debug, Clone, Default, Serialize, Deserialize, JsonSchema, PartialEq)]
#[serde(deny_unknown_fields)]
pub struct BackfillSpec {
#[serde(default, skip_serializing_if = "Option::is_none")]
pub window: Option<String>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub concurrency: Option<usize>,
#[serde(default, skip_serializing_if = "Option::is_none")]
pub timezone: Option<String>,
}
impl BackfillSpec {
pub fn validate(&self, source_configs: &[String]) -> CliResult<()> {
if let Some(w) = &self.window {
crate::backfill::plan::parse_window(w)?;
}
if self.concurrency == Some(0) {
return Err(CliError::Config(
"backfill.concurrency must be at least 1".into(),
));
}
if let Some(tz) = &self.timezone {
parse_timezone(tz)?;
}
if !source_configs.is_empty() && !source_configs.iter().any(|c| has_scoping_tokens(c)) {
return Err(CliError::Config(
"the config has a `backfill:` block but no source config references a \
`${backfill.start}` / `${backfill.end}` / `${now.*}` token — every window \
would replay identical data. Scope the source to the window (e.g. \
`query: SELECT * FROM t WHERE updated_at >= '${backfill.start}' AND \
updated_at < '${backfill.end}'`), or drop the block and use \
`faucet backfill --from-bookmark` instead"
.into(),
));
}
Ok(())
}
}
pub fn has_scoping_tokens(serialized_config: &str) -> bool {
serialized_config.contains("${backfill.") || serialized_config.contains("${now.")
}
pub fn parse_timezone(name: &str) -> CliResult<chrono_tz::Tz> {
name.parse::<chrono_tz::Tz>().map_err(|_| {
CliError::Config(format!(
"'{name}' is not a valid IANA timezone (e.g. UTC, America/New_York)"
))
})
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn parses_full_block() {
let yaml = "window: 1d\nconcurrency: 4\ntimezone: America/New_York\n";
let spec: BackfillSpec = serde_yaml::from_str(yaml).unwrap();
assert_eq!(spec.window.as_deref(), Some("1d"));
assert_eq!(spec.concurrency, Some(4));
spec.validate(&["${backfill.start}".into()]).unwrap();
}
#[test]
fn rejects_unknown_field() {
assert!(serde_yaml::from_str::<BackfillSpec>("bogus: 1\n").is_err());
}
#[test]
fn rejects_bad_window_concurrency_timezone() {
let spec = BackfillSpec {
window: Some("soon".into()),
..Default::default()
};
assert!(spec.validate(&[]).is_err());
let spec = BackfillSpec {
concurrency: Some(0),
..Default::default()
};
assert!(spec.validate(&[]).is_err());
let spec = BackfillSpec {
timezone: Some("Mars/Olympus".into()),
..Default::default()
};
assert!(spec.validate(&[]).is_err());
}
#[test]
fn rejects_unscoped_source_with_block() {
let spec = BackfillSpec::default();
let err = spec
.validate(&[r#"{"query":"SELECT * FROM t"}"#.into()])
.unwrap_err();
assert!(err.to_string().contains("${backfill.start}"), "{err}");
spec.validate(&[r#"{"prefix":"dt=${now.date}/"}"#.into()])
.unwrap();
spec.validate(&[]).unwrap();
}
#[test]
fn timezone_parses() {
parse_timezone("UTC").unwrap();
parse_timezone("America/New_York").unwrap();
assert!(parse_timezone("Nowhere").is_err());
}
}