pub struct RestStreamConfig {Show 41 fields
pub base_url: String,
pub path: String,
pub method: Method,
pub auth: AuthSpec<Auth>,
pub headers: HashMap<String, String>,
pub query_params: HashMap<String, String>,
pub query_params_multi: HashMap<String, Vec<String>>,
pub body: Option<Value>,
pub pagination: PaginationStyle,
pub records_path: Option<String>,
pub max_pages: Option<usize>,
pub request_delay: Option<Duration>,
pub timeout: Option<Duration>,
pub max_retries: u32,
pub retry_backoff: Duration,
pub tolerated_http_errors: Vec<u16>,
pub replication_method: ReplicationMethod,
pub replication_key: Option<String>,
pub start_replication_value: Option<Value>,
pub state_key: Option<String>,
pub name: Option<String>,
pub primary_keys: Vec<String>,
pub schema: Option<Value>,
pub schema_sample_size: usize,
pub partitions: Vec<HashMap<String, Value>>,
pub partition_concurrency: Option<usize>,
pub tls: Option<TlsClientConfig>,
pub response_format: ResponseFormat,
pub csv_delimiter: u8,
pub csv_has_headers: bool,
pub excel_sheet: Option<String>,
pub excel_header_row: usize,
pub replication_bind: Option<ReplicationBind>,
pub odata: Option<ODataConfig>,
pub decode: Vec<DecodeStep>,
pub async_job: Option<AsyncJobConfig>,
pub window: Option<WindowSpec>,
pub record_ancestors: Option<HashMap<String, String>>,
pub records_multi: Vec<RecordsMultiSpec>,
pub op_field: Option<String>,
pub persist_cursor: bool,
}Expand description
Configuration for a RestStream.
Fields§
§base_url: String§path: StringURL path, relative to base_url. May contain {key} placeholders that
are substituted per-partition (e.g. "/orgs/{org_id}/users").
method: Method§auth: AuthSpec<Auth>Authentication: either inline ({ type, config }) or a { ref: <name> }
pointer to a shared provider in the CLI’s top-level auth: catalog.
headers: HashMap<String, String>Static request headers sent on every request (data pages, async-job
submit/poll/fetch requests, and OData $metadata discovery probes).
Applied before the auth provider’s header placements, so an auth
header of the same name always wins on a clash. Values honor
${env:} / ${param.*} load-time interpolation and pass through the
secrets/redaction boundary like other config strings. Invalid header
names/values are rejected at config load
(FaucetError::Config), never a
mid-run panic.
headers:
Prefer: transient
Accept: application/jsonquery_params: HashMap<String, String>§query_params_multi: HashMap<String, Vec<String>>Repeated / array-valued query params (#536), rendered as repeated keys —
e.g. { "group_by[]": ["api_key_id", "model"] } → ?group_by[]=api_key_id&group_by[]=model.
Applied alongside (in addition to) query_params;
use this for APIs that need a key to appear more than once (group_by[],
repeated expand/fields). Values honor {placeholder} context
substitution for child sources, like query_params. Empty by default.
body: Option<Value>§pagination: PaginationStyle§records_path: Option<String>§max_pages: Option<usize>§request_delay: Option<Duration>§timeout: Option<Duration>§max_retries: u32Number of retries (after the first attempt) for transient request
failures. Default 3.
Precedence note: the REST source predates the unified pipeline
resilience: policy. When this field (or retry_backoff)
is left at its default, an injected RetryPolicy (e.g. from a
pipeline-level resilience: block, via
RestStream::with_retry_policy)
governs the retry budget. Setting this field away from its default makes
it win — an explicit per-connector value is never silently overridden by
a pipeline-wide default.
retry_backoff: DurationBase exponential-backoff delay between retries. Default 1s. Shares the
legacy-field precedence rule documented on max_retries.
tolerated_http_errors: Vec<u16>HTTP status codes that should not cause an error. Responses with these codes are treated as empty pages (no records, no further pages).
replication_method: ReplicationMethod§replication_key: Option<String>Field name (not a JSONPath) used for incremental replication bookmarking.
start_replication_value: Option<Value>Bookmark value: records where record[replication_key] <= start_replication_value
are filtered out when replication_method is Incremental.
state_key: Option<String>Opt-in identifier used by Pipeline::with_state_store
to persist this stream’s bookmark across runs. When set, the pipeline
will load any previously-stored bookmark before fetching and write the
new bookmark only after the sink confirms the batch.
Keys must satisfy faucet_core::state::validate_state_key.
name: Option<String>Human-readable stream name (used in logging and Singer SCHEMA messages).
primary_keys: Vec<String>Field names that uniquely identify a record (Singer key_properties).
schema: Option<Value>JSON Schema describing the structure of each record.
schema_sample_size: usizeMaximum number of records to sample when inferring the schema via
crate::stream::RestStream::infer_schema. 0 means sample all
available records (up to max_pages). Defaults to 100.
partitions: Vec<HashMap<String, Value>>Each entry is a context map whose values are substituted into path
placeholders. The stream is executed once per partition and results are
concatenated. Empty means run once with no substitution.
partition_concurrency: Option<usize>Maximum number of partitions to fetch concurrently.
None means sequential processing (backward compatible default).
tls: Option<TlsClientConfig>Optional client-certificate (mutual TLS) config. When set, the source
presents a client certificate on every request — data requests and
any inline auth token request (both go through the same HTTP client).
Requires the crate’s mtls feature; a tls block on a build without it
is a load-time error rather than being silently ignored.
response_format: ResponseFormatHow to parse the response body. json (default) uses JSONPath
extraction (records_path); csv / excel parse a tabular file
body into records — for authenticated file endpoints such as a Microsoft
Graph / OneDrive / SharePoint …/content download or any signed export
URL. In file mode a single response is fetched (pagination must be
none) and records_path does not apply. excel requires the crate’s
excel feature.
csv_delimiter: u8CSV field delimiter byte (default ,). Used only when
response_format: csv.
csv_has_headers: boolWhether the first CSV row is a header row supplying field names
(default true). When false, fields are named column_0, column_1, …
excel_sheet: Option<String>Excel worksheet to read: a sheet name, or a 0-based index as a string.
When omitted, the first worksheet is used. response_format: excel only.
excel_header_row: usize0-based index of the Excel header row (default 0). Rows above it are
skipped; the header row supplies field names. response_format: excel only.
replication_bind: Option<ReplicationBind>Bind the stored bookmark into the outgoing request (query param / header /
body field / path) so the server returns only new rows. Composes with the
existing replication_key client-side filter, which stays active as a
safety net. Requires replication_method: incremental + replication_key.
odata: Option<ODataConfig>Speak the OData protocol: @odata.nextLink paging, the $.value
envelope, $select/$filter/$expand/$orderby sugar, and
$metadata (EDMX) → schema discovery. When set, it derives the
pagination, records_path, query params, and Prefer header at load
time (explicit values still win). See ODataConfig.
decode: Vec<DecodeStep>Decode the response body before record extraction: a chain of
extract (JSONPath) / base64 / gunzip / unzip / parse
(json|csv|xlsx|xml) steps. Lets a source consume base64/compressed/file
payloads (e.g. a base64 XLSX inside a SOAP body, or a gzipped-CSV export).
When set, it replaces the response_format body parsing, and pagination
must be none. See DecodeStep.
async_job: Option<AsyncJobConfig>Run a submit→poll→fetch job lifecycle instead of a single GET, for
bulk/export/report-run APIs (Salesforce Bulk, Stripe Reporting, …). The
fetched result flows through decode: / response_format. When set,
pagination must be none. See AsyncJobConfig.
window: Option<WindowSpec>Bound each request to a rolling [start, end) window between the stored
bookmark and now, iterating the windows within one run (each step
wide) with per-window bookmark durability. For APIs that require — or cap —
a bounded date range (analytics/ads/reporting feeds). Parity with Airbyte’s
DatetimeBasedCursor. Requires replication_method: incremental +
replication_key, and a start bookmark (from state, or
start_replication_value). See WindowSpec.
record_ancestors: Option<HashMap<String, String>>When records_path selects a nested array element (e.g.
$.data[*].data.object), copy fields from the enclosing [*]
array-element ancestor onto each emitted record. The map is
dest_field: ancestor_relative_path — for each matched leaf, the source
walks up to the array-element ancestor and copies the named path onto the
record under dest_field. Absent ⇒ records are emitted unchanged.
records_path: "$.data[*].data.object"
record_ancestors: { event_id: "id", event_created: "created" }Requires records_path to contain an array wildcard [*]; mutually
exclusive with records_multi.
records_multi: Vec<RecordsMultiSpec>Emit several record arrays from one response in a single page (sharing one
pagination advance), each stamped with a user-defined op marker under
op_field. Composes with a downstream
write_mode: upsert + delete_marker so added/modified/removed feeds
route correctly. Mutually exclusive with records_path /
record_ancestors, and requires response_format: json with no
decode: pipeline.
records_multi:
- { path: "$.added[*]", op: upsert }
- { path: "$.modified[*]", op: upsert }
- { path: "$.removed[*]", op: delete }
op_field: _opop_field: Option<String>Field name each records_multi record is stamped
with its spec’s op value. Defaults to _op when omitted.
persist_cursor: boolPersist the terminal pagination cursor as this run’s bookmark (riding the
existing StreamPage.bookmark / StateStore path — no core trait change)
and, on resume, seed the stored bookmark back into the first request
(query param for cursor, request body field for cursor_in_body) before
paging. Only meaningful with pagination: cursor / cursor_in_body;
mutually exclusive with window slicing. Default false.
Implementations§
Source§impl RestStreamConfig
impl RestStreamConfig
Sourcepub fn validate(&self) -> Result<(), FaucetError>
pub fn validate(&self) -> Result<(), FaucetError>
Validate cross-field invariants that serde alone can’t express.
File response formats (csv / excel) fetch a single response and
parse the whole body, so paginated / JSONPath-extracted requests are
rejected rather than silently ignored.
Sourcepub fn apply_odata_defaults(&mut self)
pub fn apply_odata_defaults(&mut self)
Derive request defaults from the odata: block (paging, $.value
envelope, $select/$filter/$expand/$orderby params, and the
Prefer page-size header). Explicit config always wins — a field the
user already set is never overwritten. Idempotent.
pub fn new(base_url: &str, path: &str) -> Self
pub fn method(self, m: Method) -> Self
pub fn auth(self, a: Auth) -> Self
Sourcepub fn header(self, k: &str, v: &str) -> Self
pub fn header(self, k: &str, v: &str) -> Self
Add a static request header. Validation is deferred to
RestStream::new (via validate),
so an invalid name/value surfaces as a typed
FaucetError::Config rather than
panicking here.
pub fn query(self, k: &str, v: &str) -> Self
pub fn body(self, b: Value) -> Self
Sourcepub fn tls(self, tls: TlsClientConfig) -> Self
pub fn tls(self, tls: TlsClientConfig) -> Self
Attach a mutual-TLS client identity (requires the mtls feature at build
time; otherwise RestStream::new errors).
pub fn pagination(self, p: PaginationStyle) -> Self
pub fn records_path(self, p: &str) -> Self
pub fn max_pages(self, n: usize) -> Self
pub fn request_delay(self, d: Duration) -> Self
pub fn timeout(self, d: Duration) -> Self
pub fn max_retries(self, n: u32) -> Self
pub fn retry_backoff(self, d: Duration) -> Self
Sourcepub fn tolerate_http_error(self, status: u16) -> Self
pub fn tolerate_http_error(self, status: u16) -> Self
HTTP status codes that should be silently ignored (treated as empty pages).
pub fn replication_method(self, m: ReplicationMethod) -> Self
Sourcepub fn replication_key(self, key: &str) -> Self
pub fn replication_key(self, key: &str) -> Self
Field name (not JSONPath) used as the incremental replication bookmark.
Sourcepub fn start_replication_value(self, v: Value) -> Self
pub fn start_replication_value(self, v: Value) -> Self
Bookmark start value: records at or before this value are filtered out
when using ReplicationMethod::Incremental.
Sourcepub fn state_key(self, key: &str) -> Self
pub fn state_key(self, key: &str) -> Self
Opt the stream into resumable runs by giving it a stable state key.
When this is set and the Pipeline is
configured with a state store, the previously persisted bookmark is
applied to the stream before fetching.
Sourcepub fn replication_bind(self, bind: ReplicationBind) -> Self
pub fn replication_bind(self, bind: ReplicationBind) -> Self
Bind the stored bookmark into the outgoing request (#513).
Sourcepub fn window(self, window: WindowSpec) -> Self
pub fn window(self, window: WindowSpec) -> Self
Slice the run into rolling [start, end) datetime windows (#527).
Sourcepub fn odata(self, odata: ODataConfig) -> Self
pub fn odata(self, odata: ODataConfig) -> Self
Speak OData: derive paging, the $.value envelope, the query-option
sugar, and $metadata discovery from the block (#512).
Sourcepub fn decode(self, steps: Vec<DecodeStep>) -> Self
pub fn decode(self, steps: Vec<DecodeStep>) -> Self
Set the response-decode pipeline (#515).
Sourcepub fn primary_keys(self, keys: Vec<String>) -> Self
pub fn primary_keys(self, keys: Vec<String>) -> Self
Field names that uniquely identify a record (Singer key_properties).
Sourcepub fn schema_sample_size(self, n: usize) -> Self
pub fn schema_sample_size(self, n: usize) -> Self
Maximum records to sample for schema inference (0 = unlimited).
Sourcepub fn add_partition(self, ctx: HashMap<String, Value>) -> Self
pub fn add_partition(self, ctx: HashMap<String, Value>) -> Self
Add a partition context. The stream will execute once for each partition,
substituting {key} placeholders in path with values from the context.
Sourcepub fn add_query_param_multi(self, key: &str, values: Vec<String>) -> Self
pub fn add_query_param_multi(self, key: &str, values: Vec<String>) -> Self
Add a repeated / array-valued query parameter (#536): key is emitted
once per value (?key=v0&key=v1). Chainable.
Sourcepub fn partition_concurrency(self, concurrency: Option<usize>) -> Self
pub fn partition_concurrency(self, concurrency: Option<usize>) -> Self
Set the maximum number of partitions to fetch concurrently.
None (default) means sequential processing.
Sourcepub fn record_ancestors(self, map: HashMap<String, String>) -> Self
pub fn record_ancestors(self, map: HashMap<String, String>) -> Self
Copy enclosing [*] ancestor fields onto each nested record (#549).
Sourcepub fn records_multi(self, specs: Vec<RecordsMultiSpec>) -> Self
pub fn records_multi(self, specs: Vec<RecordsMultiSpec>) -> Self
Emit several op-stamped record arrays from one response in one page (#548).
Sourcepub fn op_field(self, field: &str) -> Self
pub fn op_field(self, field: &str) -> Self
Field name each records_multi record is stamped
with its op value (default _op).
Sourcepub fn persist_cursor(self, enabled: bool) -> Self
pub fn persist_cursor(self, enabled: bool) -> Self
Persist the terminal pagination cursor as the run’s bookmark and seed it on resume (#547).
Trait Implementations§
Source§impl Clone for RestStreamConfig
impl Clone for RestStreamConfig
Source§fn clone(&self) -> RestStreamConfig
fn clone(&self) -> RestStreamConfig
1.0.0 (const: unstable) · Source§fn clone_from(&mut self, source: &Self)
fn clone_from(&mut self, source: &Self)
source. Read moreSource§impl Debug for RestStreamConfig
impl Debug for RestStreamConfig
Source§impl Default for RestStreamConfig
impl Default for RestStreamConfig
Source§impl<'de> Deserialize<'de> for RestStreamConfig
impl<'de> Deserialize<'de> for RestStreamConfig
Source§fn deserialize<__D>(__deserializer: __D) -> Result<Self, __D::Error>where
__D: Deserializer<'de>,
fn deserialize<__D>(__deserializer: __D) -> Result<Self, __D::Error>where
__D: Deserializer<'de>,
Source§impl JsonSchema for RestStreamConfig
impl JsonSchema for RestStreamConfig
Source§fn schema_id() -> Cow<'static, str>
fn schema_id() -> Cow<'static, str>
Source§fn json_schema(generator: &mut SchemaGenerator) -> Schema
fn json_schema(generator: &mut SchemaGenerator) -> Schema
Source§fn inline_schema() -> bool
fn inline_schema() -> bool
$ref keyword. Read more