use crate::config::PipelineConfig;
use crate::expand::{ExpandedNode, expand};
use crate::serve::error::ServeError;
use serde::{Deserialize, Serialize};
use serde_json::Value;
const SUBMIT_DISCOVERY_TIMEOUT: std::time::Duration = std::time::Duration::from_secs(60);
#[derive(Debug, Clone, Copy, PartialEq, Eq, Default, Serialize, Deserialize)]
#[serde(rename_all = "lowercase")]
pub enum ConfigFormat {
#[default]
Yaml,
Json,
}
#[derive(Debug)]
pub struct LoadedSubmission {
pub cfg: PipelineConfig,
pub nodes: Vec<ExpandedNode>,
pub tenant: Option<std::sync::Arc<TenantScope>>,
}
impl LoadedSubmission {
pub fn auth_catalog(&self) -> crate::error::CliResult<crate::auth_catalog::AuthCatalog> {
build_catalog(self.tenant.as_deref(), &self.cfg)
}
}
fn build_catalog(
tenant: Option<&TenantScope>,
cfg: &PipelineConfig,
) -> crate::error::CliResult<crate::auth_catalog::AuthCatalog> {
match tenant {
Some(t) => (t.build_catalog)(cfg.auth.as_ref()),
None => crate::auth_catalog::build_auth_catalog(cfg.auth.as_ref()),
}
}
pub type CatalogBuilder = std::sync::Arc<
dyn Fn(
Option<&std::collections::HashMap<String, Value>>,
) -> crate::error::CliResult<crate::auth_catalog::AuthCatalog>
+ Send
+ Sync,
>;
pub struct TenantScope {
pub values: crate::tenant_tokens::TenantValues,
pub connections: std::collections::BTreeMap<String, Value>,
pub blocked: std::collections::BTreeMap<String, String>,
pub budget: Option<faucet_core::BudgetSpec>,
pub build_catalog: CatalogBuilder,
pub on_state_key: Option<crate::executor::StateKeyHook>,
}
impl std::fmt::Debug for TenantScope {
fn fmt(&self, f: &mut std::fmt::Formatter<'_>) -> std::fmt::Result {
f.debug_struct("TenantScope")
.field("tenant", &self.values.id)
.field("connections", &self.connections.keys().collect::<Vec<_>>())
.field("blocked", &self.blocked)
.finish()
}
}
impl TenantScope {
pub fn state_scope(&self) -> crate::executor::StateScope {
crate::executor::StateScope {
namespace: Some(self.values.id.clone()),
on_key: self.on_state_key.clone(),
}
}
fn check_refs(&self, cfg: &PipelineConfig, nodes: &[ExpandedNode]) -> Result<(), ServeError> {
for node in nodes {
for config in [&node.source.config, &node.sink.config] {
let Some(name) = crate::auth_catalog::auth_ref(config) else {
continue;
};
if let Some(why) = self.blocked.get(&name) {
return Err(ServeError::Conflict(format!(
"connection '{name}' of tenant '{}' needs re-authorization ({why}); \
reconnect it before running",
self.values.id
)));
}
if !cfg.auth.as_ref().is_some_and(|a| a.contains_key(&name)) {
return Err(ServeError::Conflict(format!(
"tenant '{}' has no connection '{name}' (row '{}' references it); \
create it with POST /v1/tenants/{}/connections or a connect flow",
self.values.id, node.id, self.values.id
)));
}
}
}
Ok(())
}
}
#[cfg(feature = "policy")]
pub type ServerPolicy = faucet_core::PolicySpec;
#[cfg(not(feature = "policy"))]
pub type ServerPolicy = ();
pub async fn load_submission(
body: &str,
format: ConfigFormat,
default_base: Option<&Value>,
policy: Option<&ServerPolicy>,
) -> Result<LoadedSubmission, ServeError> {
load_submission_scoped(body, format, default_base, policy, None).await
}
pub async fn load_submission_scoped(
body: &str,
format: ConfigFormat,
default_base: Option<&Value>,
policy: Option<&ServerPolicy>,
tenant: Option<std::sync::Arc<TenantScope>>,
) -> Result<LoadedSubmission, ServeError> {
let mut submitted: Value = match format {
ConfigFormat::Yaml => serde_yaml::from_str(body)
.map_err(|e| ServeError::BadConfig(format!("invalid YAML: {e}")))?,
ConfigFormat::Json => serde_json::from_str(body)
.map_err(|e| ServeError::BadConfig(format!("invalid JSON: {e}")))?,
};
crate::interpolate::interpolate_value(&mut submitted)
.map_err(|e| ServeError::BadConfig(e.to_string()))?;
let mut merged = match default_base {
Some(base) => {
let mut m = base.clone();
crate::merge::merge_value(&mut m, submitted);
m
}
None => submitted,
};
crate::tenant_tokens::bind_document(&mut merged, tenant.as_deref().map(|t| &t.values))
.map_err(|message| ServeError::Unprocessable {
message,
details: None,
})?;
crate::params::bind_document(
&mut merged,
&Default::default(),
crate::params::BindMode::Strict,
)
.map_err(|e| ServeError::Unprocessable {
message: e.to_string(),
details: None,
})?;
let mut cfg = PipelineConfig::from_value(merged).map_err(|e| ServeError::Unprocessable {
message: e.to_string(),
details: None,
})?;
#[cfg(feature = "schedule")]
if cfg.schedule.is_some() {
return Err(ServeError::BadConfig(
"submitted config contains a `schedule:` block — serve runs once per \
submission; use `faucet schedule` for cron scheduling"
.into(),
));
}
#[cfg(feature = "policy")]
if let Some(server_policy) = policy {
cfg.policy = Some(match cfg.policy.take() {
Some(own) => own
.merge(server_policy.clone())
.map_err(|e| ServeError::BadConfig(format!("policy: {e}")))?,
None => server_policy.clone(),
});
}
#[cfg(not(feature = "policy"))]
let _ = policy;
crate::secrets::resolve_secrets(&mut cfg)
.await
.map_err(|e| ServeError::BadConfig(e.to_string()))?;
if let Some(t) = &tenant {
let catalog = cfg.auth.get_or_insert_with(Default::default);
for (name, spec) in &t.connections {
catalog.insert(name.clone(), spec.clone());
}
}
let auth =
build_catalog(tenant.as_deref(), &cfg).map_err(|e| ServeError::BadConfig(e.to_string()))?;
tokio::time::timeout(
SUBMIT_DISCOVERY_TIMEOUT,
crate::dynamic_fanout::resolve_dynamic_fanout(&mut cfg, &auth),
)
.await
.map_err(|_| ServeError::Unprocessable {
message: format!(
"discovery fan-out did not complete within {}s — the source's describe \
endpoint is slow or unreachable",
SUBMIT_DISCOVERY_TIMEOUT.as_secs()
),
details: None,
})?
.map_err(|e| ServeError::Unprocessable {
message: e.to_string(),
details: None,
})?;
let nodes = expand(&cfg).map_err(|e| ServeError::Unprocessable {
message: e.to_string(),
details: None,
})?;
if let Some(t) = &tenant {
t.check_refs(&cfg, &nodes)?;
}
Ok(LoadedSubmission { cfg, nodes, tenant })
}
#[cfg(test)]
mod tests {
use super::*;
use serde_json::json;
fn base() -> Value {
json!({
"version": 1,
"pipeline": {
"source": { "type": "csv", "config": { "path": "DEFAULT.csv" } },
"sink": { "type": "jsonl", "config": { "path": "out.jsonl" } }
}
})
}
#[tokio::test]
async fn submitted_overrides_default() {
let body = r#"{ "pipeline": { "source": { "config": { "path": "OVERRIDE.csv" } } } }"#;
let loaded = load_submission(body, ConfigFormat::Json, Some(&base()), None)
.await
.unwrap();
let node = &loaded.nodes[0];
assert_eq!(node.source.config["path"], "OVERRIDE.csv");
assert_eq!(node.sink.config["path"], "out.jsonl");
}
#[tokio::test]
async fn missing_version_without_base_is_unprocessable() {
let body = r#"{ "pipeline": {} }"#;
let err = load_submission(body, ConfigFormat::Json, None, None)
.await
.unwrap_err();
assert!(matches!(
err,
ServeError::Unprocessable { .. } | ServeError::BadConfig(_)
));
}
#[cfg(feature = "schedule")]
#[tokio::test]
async fn schedule_block_is_rejected() {
let body = r#"
version: 1
pipeline:
source: { type: csv, config: { path: x.csv } }
sink: { type: jsonl, config: { path: out.jsonl } }
schedule:
cron: "0 * * * *"
timezone: UTC
"#;
let err = load_submission(body, ConfigFormat::Yaml, None, None)
.await
.unwrap_err();
match err {
ServeError::BadConfig(m) => assert!(m.contains("schedule:")),
other => panic!("expected BadConfig, got {other:?}"),
}
}
#[tokio::test]
async fn invalid_yaml_is_bad_config() {
let err = load_submission("{[bad", ConfigFormat::Yaml, None, None)
.await
.unwrap_err();
assert!(matches!(err, ServeError::BadConfig(_)));
}
#[tokio::test]
async fn submitted_extends_is_rejected_with_composition_hint() {
let body = "version: 1\nextends: /etc/passwd\npipeline:\n source: { type: csv, config: { path: x.csv } }\n sink: { type: jsonl, config: { path: o.jsonl } }\n";
let err = load_submission(body, ConfigFormat::Yaml, None, None)
.await
.unwrap_err();
let msg = match &err {
ServeError::Unprocessable { message, .. } => message.clone(),
ServeError::BadConfig(m) => m.clone(),
other => format!("{other:?}"),
};
assert!(
msg.contains("composition"),
"submitted extends must be rejected with the composition hint, got: {msg}"
);
}
}