faucet-cli 1.13.0

Config-driven CLI runner for faucet-stream pipelines (YAML / JSON, Meltano-style)
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
65
66
67
68
69
70
71
72
73
74
75
76
77
78
79
80
81
82
83
84
85
86
87
88
89
90
91
92
93
94
95
96
97
98
99
100
101
102
103
104
105
106
107
108
109
110
111
112
113
114
115
116
117
118
119
120
121
122
123
124
125
126
127
128
129
130
131
132
133
134
135
136
137
138
139
140
141
142
143
144
145
146
147
148
149
150
151
152
153
154
155
156
157
158
159
160
161
162
163
164
165
166
167
168
169
170
171
172
173
174
175
176
177
178
179
180
181
182
183
184
185
186
187
188
189
190
191
192
193
194
195
196
197
198
199
200
201
202
203
204
205
206
207
208
209
210
211
212
213
214
215
216
217
218
219
220
221
222
223
224
225
226
227
228
229
230
231
232
233
234
235
236
237
238
239
240
241
242
243
244
245
246
247
248
249
250
251
252
253
254
255
256
257
258
259
260
261
262
263
264
265
266
267
268
269
270
271
272
273
274
275
276
277
278
279
280
281
282
283
284
285
286
287
288
289
290
291
292
293
294
295
296
297
298
299
300
301
302
303
304
305
306
307
308
309
310
311
312
313
314
315
316
317
318
319
320
321
322
323
324
325
326
327
328
329
330
331
332
333
334
335
336
337
338
339
340
341
342
343
344
345
346
347
348
349
350
351
352
353
354
355
356
357
358
359
360
361
362
363
364
365
366
367
368
369
370
371
372
373
374
375
376
377
378
379
380
381
382
383
384
385
386
387
388
389
390
391
392
393
394
395
396
397
398
399
400
401
402
403
404
405
406
407
408
409
410
411
412
413
414
415
416
417
418
#![cfg_attr(docsrs, feature(doc_cfg))]

//! # faucet-cli
//!
//! A config-driven runner for [`faucet-stream`](https://docs.rs/faucet-stream)
//! pipelines.  Define a source, optional transforms, a sink, and (optionally)
//! a state store in a YAML or JSON file, then run it with the `faucet` binary —
//! no Rust code required.
//!
//! The library half of this crate exposes the same building blocks the binary
//! uses (config parsing, env interpolation, the connector registry) so that
//! integrations and tests can reuse them.

pub mod auth_catalog;
pub mod backfill;
pub mod budget;
#[cfg(feature = "catalog")]
pub mod catalog;
pub mod chunking;
pub mod cli;
pub mod commands;
pub mod compose;
pub mod config;
pub mod conformance;
pub mod connector_export;
pub mod discovery_matrix;
pub mod dlq_replay;
pub mod dynamic_fanout;
pub mod env_config;
pub mod env_loader;
pub mod error;
pub mod exec_metrics;
pub mod executor;
pub mod expand;
pub mod hub;
#[cfg(feature = "catalog")]
pub mod impact;
pub mod init_template;
pub mod interpolate;
#[cfg(feature = "lineage")]
pub mod lineage_glue;
/// Shared live-view metrics plumbing (recorder install + Prometheus-text
/// sampler), compiled when either live-view feature is on.
#[cfg(any(feature = "cli-tui", feature = "cli-progress"))]
pub mod livemetrics;
#[cfg(feature = "serve")]
pub mod local_outputs;
#[cfg(feature = "mcp")]
pub mod mcp;
pub mod memstat;
pub mod merge;
#[cfg(feature = "notify")]
pub mod notify;
pub mod obs;
pub mod params;
pub mod partition;
pub mod pipeline_state;
pub mod pipeline_test;
#[cfg(feature = "policy")]
pub mod policy;
pub mod profiling;
#[cfg(feature = "cli-progress")]
pub mod progress;
pub mod reconcile;
pub mod registry;
pub mod registry_index;
pub mod replication;
pub mod rollback;
pub mod scaffold;
#[cfg(feature = "schedule")]
pub mod schedule;
pub mod schema_compose;
pub mod secrets;
pub mod select;
#[cfg(feature = "serve")]
pub mod serve;
pub mod signals;
pub mod sla;
pub mod state;
pub mod status;
#[cfg(feature = "templates")]
pub mod templates;
pub mod tenant_tokens;
pub mod topology;
pub mod transforms;
#[cfg(feature = "cli-tui")]
pub mod tui;
pub mod usage;
pub mod verify;
pub mod vocabulary;

pub use error::{CliError, CliResult};

use crate::cli::{Cli, Command};
use crate::registry::PluginRegistry;

/// Entry point for a custom `faucet` binary that bundles third-party
/// connectors.
///
/// A custom-CLI author writes a tiny `main.rs` that builds a [`PluginRegistry`]
/// with their connectors registered on top of the built-ins and hands it here:
///
/// ```no_run
/// use faucet_cli::registry::PluginRegistry;
/// fn main() -> std::process::ExitCode {
///     faucet_cli::run_main(PluginRegistry::with_builtins())
/// }
/// ```
///
/// This installs `registry` as the process-global connector registry (so every
/// command — `run`, `validate`, `schema`, `list`, `preview`, `serve`, … — sees
/// the custom connectors), parses argv, installs the tracing subscriber, and
/// dispatches. The return value is the process exit code (the failed-probe /
/// failed-case / failed-unit count for `doctor` / `test` / `backfill`, `1` for
/// any other error, `0` on success). The stock `faucet` binary calls this with
/// `PluginRegistry::with_builtins()`.
pub fn run_main(registry: PluginRegistry) -> std::process::ExitCode {
    use clap::Parser;
    use std::process::ExitCode;

    if let Err(err) = registry.install() {
        commands::report(&err);
        return ExitCode::from(1);
    }

    // Hand `faucet-core` the secret scrubber before anything can produce output.
    // Core builds two things that leave the process and that it cannot redact on
    // its own: the DLQ envelope's `error.message`, and error text a caller
    // forwards on. Installed here, at the single entry point, so every subcommand
    // and runtime gets it (#456 H5).
    faucet_core::redact::install(Box::new(|s: &str| {
        secrets::registry::redact(s).into_owned()
    }));

    // Dynamic shell completion (#383): when the shell invokes us with the
    // `COMPLETE` env var set, compute and print candidates, then exit — before
    // any normal parsing/tracing/runtime setup. A no-op otherwise. The registry
    // is installed above so connector-kind candidates reflect the live binary.
    clap_complete::env::CompleteEnv::with_factory(<Cli as clap::CommandFactory>::command)
        .complete();

    let cli = Cli::parse();
    #[cfg(feature = "serve")]
    let is_serve = matches!(cli.command, Command::Serve(_));
    #[cfg(not(feature = "serve"))]
    let is_serve = false;
    // A `--tui` run on a real terminal routes logs into the TUI's in-memory
    // ring (the stdout subscriber would corrupt the alternate screen).
    #[cfg(feature = "cli-tui")]
    let is_tui = matches!(&cli.command, Command::Run(a) if tui::is_tui_session(a.tui));
    #[cfg(not(feature = "cli-tui"))]
    let is_tui = false;
    // `faucet mcp` (stdio) writes JSON-RPC on stdout, so its logs must go to
    // stderr — never the default subscriber.
    #[cfg(feature = "mcp")]
    let is_mcp = matches!(cli.command, Command::Mcp(_));
    #[cfg(not(feature = "mcp"))]
    let is_mcp = false;
    // `serve` installs its own (redacting, run-scoped) subscriber; every other
    // command uses the plain redacting fmt subscriber.
    if !is_serve && !is_tui && !is_mcp {
        crate::cli::set_log_format(cli.log_format);
        install_tracing(&cli.log_level, cli.log_format);
    }
    #[cfg(feature = "cli-tui")]
    if is_tui {
        tui::install_tui_tracing(&cli.log_level);
    }
    #[cfg(feature = "mcp")]
    if is_mcp {
        crate::cli::set_log_format(cli.log_format);
        mcp::install_stderr_tracing(&cli.log_level, cli.log_format);
    }

    let runtime = match tokio::runtime::Builder::new_multi_thread()
        .enable_all()
        .build()
    {
        Ok(rt) => rt,
        Err(e) => {
            eprintln!("error: failed to start async runtime: {e}");
            return ExitCode::from(1);
        }
    };

    runtime.block_on(async move {
        match run_command(cli).await {
            Ok(()) => ExitCode::SUCCESS,
            // `doctor` / `test` / `backfill` already printed their report; the
            // exit code is the failed count (clamped to 255).
            Err(CliError::DoctorFailed { failed }) => ExitCode::from(failed.min(255) as u8),
            Err(CliError::TestsFailed { failed }) => ExitCode::from(failed.min(255) as u8),
            Err(CliError::BackfillFailed { failed }) => ExitCode::from(failed.min(255) as u8),
            // `verify` printed its report; the exit code is the differing-key
            // count. A blocked `rollback` exits with the conflict count.
            Err(CliError::PolicyViolations { violations }) => {
                ExitCode::from(violations.clamp(1, 255) as u8)
            }
            Err(CliError::VerifyFailed { differences }) => {
                ExitCode::from(differences.clamp(1, 255) as u8)
            }
            Err(CliError::RollbackBlocked { conflicts }) => {
                ExitCode::from(conflicts.clamp(1, 255) as u8)
            }
            // `status` printed its screen; 1 = degraded / unknown, 2 = failed.
            Err(CliError::StatusUnhealthy { code, .. }) => ExitCode::from(code),
            Err(err) => {
                commands::report(&err);
                ExitCode::from(1)
            }
        }
    })
}

/// Dispatch a parsed [`Cli`] to the matching command. Public so custom hosts and
/// integration tests can drive the exact same code path as [`run_main`] with a
/// programmatically-built `Cli` (and a registry installed via
/// [`PluginRegistry::install`]).
pub async fn run_command(cli: Cli) -> CliResult<()> {
    Box::pin(dispatch(cli)).await
}

async fn dispatch(cli: Cli) -> CliResult<()> {
    #[cfg(feature = "serve")]
    let serve_log_level = cli.log_level.clone();
    #[cfg(feature = "serve")]
    let log_format = cli.log_format;
    match cli.command {
        Command::Run(args) => commands::run::run(args).await,
        Command::Backfill(args) => commands::backfill::run(args).await,
        Command::Replicate(args) => commands::replicate::run(args).await,
        Command::Discover(args) => commands::discover::run(args).await,
        Command::Validate(args) => commands::validate::run(args).await,
        Command::Schema(args) => commands::schema::run(args).await,
        Command::List(args) => commands::list::run(args).await,
        Command::Search(args) => commands::search::run(args).await,
        Command::Conformance(args) => commands::conformance::run(args).await,
        Command::Install(args) => commands::install::run(args).await,
        Command::Preview(args) => commands::preview::run(args).await,
        Command::Plan(args) => commands::plan::run(args).await,
        #[cfg(feature = "cli-dev")]
        Command::Dev(args) => commands::dev::run(args).await,
        Command::Init(args) => commands::init::run(args).await,
        Command::New(args) => commands::new::run(args).await,
        Command::Doctor(args) => commands::doctor::run(args).await,
        Command::Test(args) => commands::test::run(args).await,
        Command::Dlq(args) => commands::dlq::run(args).await,
        Command::Verify(args) => commands::verify::run(args).await,
        Command::Rollback(args) => commands::rollback::run(args).await,
        Command::Profiling(args) => commands::profiling::run(args).await,
        Command::State(args) => commands::state::run(args).await,
        Command::Status(args) => commands::status::run(args).await,
        Command::Hub(args) => commands::hub::run(args).await,
        #[cfg(feature = "contract")]
        Command::Contract(args) => commands::contract::run(args).await,
        #[cfg(feature = "masking")]
        Command::Masking(args) => commands::masking::run(args).await,
        #[cfg(feature = "policy")]
        Command::Policy(args) => commands::policy::run(args).await,
        #[cfg(feature = "schedule")]
        Command::Schedule(args) => commands::schedule::run(args).await,
        #[cfg(feature = "serve")]
        Command::Serve(args) => commands::serve::run(args, serve_log_level, log_format).await,
        #[cfg(feature = "mcp")]
        Command::Mcp(args) => commands::mcp::run(args).await,
        #[cfg(feature = "notify")]
        Command::Notify(args) => commands::notify::run(args).await,
        #[cfg(feature = "catalog")]
        Command::Catalog(args) => commands::catalog::run(args).await,
        #[cfg(feature = "catalog")]
        Command::Usage(args) => commands::usage::run(args).await,
        #[cfg(feature = "templates")]
        Command::Template(args) => commands::template::run(args).await,
        Command::Completions(args) => commands::completions::run(args.shell),
        Command::Migrate(args) => commands::migrate::run(args).await,
        Command::Fmt(args) => commands::fmt::run(args).await,
        Command::Explain(args) => commands::explain::run(args).await,
        #[cfg(feature = "catalog")]
        Command::History(args) => commands::history::run(args).await,
        #[cfg(feature = "catalog")]
        Command::Cleanup(args) => commands::cleanup::run(args).await,
    }
}

#[cfg(feature = "observability")]
fn install_tracing(level: &str, format: crate::cli::LogFormat) {
    use crate::secrets::registry::RedactingMakeWriter;
    use tracing_subscriber::EnvFilter;
    let filter = EnvFilter::try_new(level).unwrap_or_else(|_| EnvFilter::new("info"));
    let builder = tracing_subscriber::fmt()
        .with_env_filter(filter)
        // Redaction wraps the *serialized bytes*, so it keeps working under
        // JSON: a resolved secret appearing in a field value is still scrubbed.
        .with_writer(RedactingMakeWriter);
    match format {
        crate::cli::LogFormat::Text => {
            let _ = builder.try_init();
        }
        crate::cli::LogFormat::Json => {
            let _ = builder
                .json()
                // Event fields at the top level rather than nested under
                // `fields`, which is what log pipelines index on.
                .flatten_event(true)
                // Span context (`pipeline`, `row`, `run_id`, …) as fields; the
                // full ancestor list is noise once the current span's fields
                // are present.
                .with_current_span(true)
                .with_span_list(false)
                .try_init();
        }
    }
}

/// Stub used when the `observability` feature is disabled. Logging falls back to
/// whatever the host environment has wired (or nothing).
#[cfg(not(feature = "observability"))]
fn install_tracing(_level: &str, _format: crate::cli::LogFormat) {}

/// Convenience entry point for integration tests and custom hosts: parse a
/// YAML config string, expand the matrix, and run all rows.
///
/// This skips the `install_observability` call that [`commands::run::run`]
/// performs — callers can wire their own `metrics` recorder / tracing
/// subscriber before calling this function (or not at all).
pub async fn run_from_yaml_str(yaml: &str) -> CliResult<executor::RunSummary> {
    run_from_yaml_str_selected(yaml, None).await
}

/// [`run_from_yaml_str`] over a subset of the config's matrix rows (#741).
pub async fn run_from_yaml_str_selected(
    yaml: &str,
    selection: Option<&select::SelectionRequest>,
) -> CliResult<executor::RunSummary> {
    // Parse first, then resolve ${env}/${file}/${secret} INTO the parsed tree
    // (post-parse) so a resolved value can never alter the document's structure
    // (F43) — mirroring the binary's `from_path` path.
    let mut value: serde_json::Value =
        serde_yaml::from_str(yaml).map_err(|e| CliError::ParseConfig {
            path: std::path::PathBuf::from("<yaml-string>"),
            message: e.to_string(),
        })?;
    interpolate::interpolate_value(&mut value)?;
    // Bind `${param.*}` from the config's own `params:` defaults (#444). No
    // caller-supplied values on this convenience path, so a required param is a
    // clear error rather than a token leaking into a connector config.
    params::bind_document(&mut value, &Default::default(), params::BindMode::Strict)?;
    let interpolated = serde_yaml::to_string(&value).map_err(|e| CliError::ParseConfig {
        path: std::path::PathBuf::from("<yaml-string>"),
        message: e.to_string(),
    })?;
    let mut cfg: config::PipelineConfig =
        serde_yaml::from_str(&interpolated).map_err(|e| CliError::ParseConfig {
            path: std::path::PathBuf::from("<yaml-string>"),
            message: e.to_string(),
        })?;
    if cfg.version != 1 {
        return Err(CliError::ParseConfig {
            path: std::path::PathBuf::from("<yaml-string>"),
            message: format!(
                "unsupported pipeline version {}, only version 1 is recognised",
                cfg.version
            ),
        });
    }
    crate::secrets::resolve_secrets(&mut cfg).await?;
    let pipeline_name = cfg.name.clone().unwrap_or_else(|| "unnamed".to_string());
    let auth = auth_catalog::build_auth_catalog(cfg.auth.as_ref())?;
    let resilience = match &cfg.resilience {
        Some(spec) => Some(spec.to_policy()?),
        None => None,
    };
    #[cfg(feature = "catalog")]
    let catalog = match cfg.catalog.as_ref() {
        Some(spec) => Some(catalog::connect_from_spec(spec).await?),
        None => None,
    };
    let nodes = expand::expand(&cfg)?;
    let nodes = match selection {
        Some(sel) => sel.apply(&cfg, nodes)?,
        None => nodes,
    };
    executor::run_expanded(
        nodes,
        executor::ExecuteOptions {
            legacy_state_writes: false,
            pipeline_name,
            run_id: None,
            execution: cfg.execution.clone(),
            concurrency: None,
            dry_run: false,
            limit: None,
            state_path_override: None,
            state_scope: Default::default(),
            shard: None,
            auth,
            clock: chrono::Utc::now().fixed_offset(),
            cancel: None,
            resilience,
            sla: cfg.sla.clone(),
            reconcile: cfg.reconcile.clone(),
            verify: cfg.verify.clone(),
            rollback: cfg.rollback.clone(),
            #[cfg(feature = "lineage")]
            lineage: None,
            #[cfg(feature = "lineage")]
            lineage_cfg: None,
            #[cfg(feature = "notify")]
            notifier: None,
            #[cfg(feature = "catalog")]
            catalog,
            usage: usage::UsageOptions::from_spec(cfg.usage.as_ref(), None)
                .map_err(CliError::Config)?,
            budget: cfg.budget.clone(),
        },
    )
    .await
}