Skip to main content

faucet_sink_http/
sink.rs

1//! HTTP sink executor.
2
3use crate::config::{HttpBatchMode, HttpSinkAuth, HttpSinkConfig};
4use async_trait::async_trait;
5use faucet_core::util::{DEFAULT_ERROR_BODY_MAX_LEN, check_http_response};
6use faucet_core::{AuthSpec, Credential, FaucetError, SharedAuthProvider};
7use futures::stream::{FuturesUnordered, StreamExt};
8use serde_json::Value;
9use std::collections::HashMap;
10
11/// Map a [`Credential`] from a shared provider onto the [`HttpSinkAuth`]
12/// representation so the existing header-application path can be reused.
13fn credential_to_auth(cred: Credential) -> HttpSinkAuth {
14    match cred {
15        Credential::Bearer(token) => HttpSinkAuth::Bearer { token },
16        Credential::Token(token) => HttpSinkAuth::Custom {
17            headers: HashMap::from([("Authorization".to_string(), token)]),
18        },
19        Credential::Basic { username, password } => HttpSinkAuth::Basic { username, password },
20        Credential::Header { name, value } => HttpSinkAuth::Custom {
21            headers: HashMap::from([(name, value)]),
22        },
23    }
24}
25
26/// An HTTP sink that sends records to an HTTP endpoint.
27pub struct HttpSink {
28    config: HttpSinkConfig,
29    client: reqwest::Client,
30    /// Optional shared auth provider. When set, it takes precedence over inline
31    /// auth. Set via [`HttpSink::with_auth_provider`].
32    auth_provider: Option<SharedAuthProvider>,
33}
34
35impl HttpSink {
36    /// Create a new HTTP sink from the given configuration.
37    pub fn new(config: HttpSinkConfig) -> Self {
38        Self {
39            config,
40            client: reqwest::Client::new(),
41            auth_provider: None,
42        }
43    }
44
45    /// Attach a shared [`AuthProvider`](faucet_core::AuthProvider). When set,
46    /// the provider supplies the credential for every request (taking
47    /// precedence over inline auth), so several sinks can share one token with
48    /// single-flight refresh. Used by the CLI to resolve `auth: { ref }`, and
49    /// by library callers who construct one provider and inject it into many
50    /// sinks.
51    pub fn with_auth_provider(mut self, provider: SharedAuthProvider) -> Self {
52        self.auth_provider = Some(provider);
53        self
54    }
55
56    /// Resolve the effective auth for the current batch. The provider (if any)
57    /// takes precedence; otherwise inline auth is used. A bare
58    /// `AuthSpec::Reference` with no provider is an error.
59    async fn resolve_auth(&self) -> Result<HttpSinkAuth, FaucetError> {
60        if let Some(provider) = &self.auth_provider {
61            Ok(credential_to_auth(provider.credential().await?))
62        } else {
63            match &self.config.auth {
64                AuthSpec::Inline(a) => Ok(a.clone()),
65                AuthSpec::Reference(r) => Err(FaucetError::Auth(format!(
66                    "auth references provider '{}' but no provider was supplied; \
67                     set one via the CLI `auth:` catalog or `with_auth_provider`",
68                    r.name
69                ))),
70            }
71        }
72    }
73
74    /// Build an HTTP request with auth and headers applied.
75    fn apply_auth(
76        &self,
77        mut req: reqwest::RequestBuilder,
78        auth: &HttpSinkAuth,
79    ) -> Result<reqwest::RequestBuilder, FaucetError> {
80        match auth {
81            HttpSinkAuth::None => {}
82            HttpSinkAuth::Bearer { token } => {
83                req = req.bearer_auth(token);
84            }
85            HttpSinkAuth::Basic { username, password } => {
86                req = req.basic_auth(username, Some(password));
87            }
88            HttpSinkAuth::Custom { headers } => {
89                let mut hm = reqwest::header::HeaderMap::new();
90                for (name, value) in headers {
91                    let n =
92                        reqwest::header::HeaderName::from_bytes(name.as_bytes()).map_err(|e| {
93                            FaucetError::Auth(format!("invalid custom header name {name:?}: {e}"))
94                        })?;
95                    let v = reqwest::header::HeaderValue::from_str(value).map_err(|e| {
96                        FaucetError::Auth(format!("invalid custom header value for {name:?}: {e}"))
97                    })?;
98                    hm.insert(n, v);
99                }
100                req = req.headers(hm);
101            }
102        }
103        Ok(req)
104    }
105
106    /// Build an HTTP request with the given pre-resolved auth and body.
107    fn build_request_with_auth(
108        &self,
109        body: &Value,
110        auth: &HttpSinkAuth,
111    ) -> Result<reqwest::RequestBuilder, FaucetError> {
112        let req = self
113            .client
114            .request(self.config.method.clone(), &self.config.url)
115            .headers(self.config.headers.clone())
116            .json(body);
117        self.apply_auth(req, auth)
118    }
119
120    /// Send a single request with retry logic, using the pre-resolved `auth`.
121    async fn send_with_retry(&self, body: &Value, auth: &HttpSinkAuth) -> Result<(), FaucetError> {
122        let mut last_error = None;
123
124        for attempt in 0..=self.config.max_retries {
125            let req = self.build_request_with_auth(body, auth)?;
126
127            match req.send().await {
128                Ok(resp) => match check_http_response(resp, DEFAULT_ERROR_BODY_MAX_LEN).await {
129                    Ok(_) => return Ok(()),
130                    Err(e) => {
131                        if attempt < self.config.max_retries && e.is_retriable() {
132                            tracing::warn!(
133                                attempt = attempt + 1,
134                                max_retries = self.config.max_retries,
135                                error = %e,
136                                "retrying request"
137                            );
138                            last_error = Some(e);
139                            continue;
140                        }
141                        return Err(e);
142                    }
143                },
144                Err(e) => {
145                    let faucet_err = FaucetError::Http(e);
146                    if attempt < self.config.max_retries && faucet_err.is_retriable() {
147                        tracing::warn!(
148                            attempt = attempt + 1,
149                            max_retries = self.config.max_retries,
150                            error = %faucet_err,
151                            "retrying request"
152                        );
153                        last_error = Some(faucet_err);
154                        continue;
155                    }
156                    return Err(faucet_err);
157                }
158            }
159        }
160
161        Err(last_error.unwrap_or_else(|| FaucetError::Sink("max retries exhausted".into())))
162    }
163}
164
165#[async_trait]
166impl faucet_core::Sink for HttpSink {
167    fn connector_name(&self) -> &'static str {
168        "http"
169    }
170
171    fn config_schema(&self) -> serde_json::Value {
172        serde_json::to_value(faucet_core::schema_for!(HttpSinkConfig))
173            .expect("schema serialization")
174    }
175
176    fn dataset_uri(&self) -> String {
177        faucet_core::redact_uri_credentials(&self.config.url)
178    }
179
180    /// Non-mutating preflight probe (probe name `"network"`).
181    ///
182    /// Issues a lightweight `HEAD` request to the configured endpoint over the
183    /// existing reqwest client. We only care that the host is reachable — that
184    /// DNS, TCP, TLS and the server all work — so **any** HTTP response (2xx,
185    /// 4xx including `405 Method Not Allowed`, or 5xx) counts as a pass. Only a
186    /// transport/connection error (no response at all) is a failure.
187    async fn check(
188        &self,
189        ctx: &faucet_core::check::CheckContext,
190    ) -> Result<faucet_core::check::CheckReport, FaucetError> {
191        use faucet_core::check::{CheckReport, Probe};
192
193        // Resolve auth so authenticated endpoints don't reject the connection
194        // before we learn the host is reachable. An unresolvable auth ref is a
195        // configuration failure surfaced on this probe.
196        let auth = match self.resolve_auth().await {
197            Ok(a) => a,
198            Err(e) => {
199                return Ok(CheckReport::single(Probe::fail_hint(
200                    "network",
201                    std::time::Duration::ZERO,
202                    e.to_string(),
203                    "check the configured auth / that a shared auth provider is wired up",
204                )));
205            }
206        };
207
208        let started = std::time::Instant::now();
209        let hint = "check the url / DNS / TLS / that the host is reachable";
210
211        let req = self
212            .client
213            .head(&self.config.url)
214            .headers(self.config.headers.clone());
215        let req = match self.apply_auth(req, &auth) {
216            Ok(r) => r,
217            Err(e) => {
218                return Ok(CheckReport::single(Probe::fail_hint(
219                    "network",
220                    started.elapsed(),
221                    e.to_string(),
222                    hint,
223                )));
224            }
225        };
226
227        let probe = match tokio::time::timeout(ctx.timeout, req.send()).await {
228            // Any HTTP response means DNS + TCP + TLS + the host all work.
229            Ok(Ok(_)) => Probe::pass("network", started.elapsed()),
230            // Transport/connection error: no response received.
231            Ok(Err(e)) => Probe::fail_hint("network", started.elapsed(), e.to_string(), hint),
232            Err(_) => Probe::fail_hint("network", started.elapsed(), "timed out", hint),
233        };
234        Ok(CheckReport::single(probe))
235    }
236
237    async fn write_batch(&self, records: &[Value]) -> Result<usize, FaucetError> {
238        if records.is_empty() {
239            return Ok(0);
240        }
241
242        // Resolve auth once per batch (provider-first, then inline).
243        let auth = self.resolve_auth().await?;
244
245        match &self.config.batch_mode {
246            HttpBatchMode::Individual => {
247                // Run `send_with_retry` for every record with at most
248                // `concurrency` in-flight at once. We drive a
249                // `FuturesUnordered` directly, refilling it as each future
250                // completes, instead of acquiring permits up-front the way
251                // the previous semaphore-based code did — that approach
252                // deadlocked because permits were acquired sequentially in
253                // a loop before any future actually ran, so after
254                // `concurrency` iterations the next `acquire_owned().await`
255                // would block forever (closes #59).
256                let concurrency = self.config.concurrency.max(1);
257                let mut in_flight = FuturesUnordered::new();
258                let mut iter = records.iter();
259                for record in iter.by_ref().take(concurrency) {
260                    in_flight.push(self.send_with_retry(record, &auth));
261                }
262                while let Some(result) = in_flight.next().await {
263                    result?;
264                    if let Some(record) = iter.next() {
265                        in_flight.push(self.send_with_retry(record, &auth));
266                    }
267                }
268
269                tracing::debug!(records = records.len(), "HTTP individual batch written");
270                Ok(records.len())
271            }
272            HttpBatchMode::Array => {
273                // `batch_size = 0` is the "no batching" sentinel: forward
274                // whatever upstream handed us as a single JSON-array POST,
275                // preserving `StreamPage` framing. Otherwise re-chunk into
276                // `batch_size` slices and issue one POST per chunk.
277                let effective_chunk = if self.config.batch_size == 0 {
278                    records.len()
279                } else {
280                    self.config.batch_size
281                };
282
283                let mut total = 0;
284                for chunk in records.chunks(effective_chunk) {
285                    let array = Value::Array(chunk.to_vec());
286                    self.send_with_retry(&array, &auth).await?;
287                    total += chunk.len();
288                }
289                tracing::debug!(
290                    records = total,
291                    batch_size = self.config.batch_size,
292                    "HTTP array batch written"
293                );
294                Ok(total)
295            }
296        }
297    }
298
299    /// Report per-row outcomes so the DLQ router dead-letters only the records
300    /// that genuinely failed.
301    ///
302    /// In **Individual** mode every record is an independent POST, so each
303    /// record's success/failure is attributable: we attempt *all* of them
304    /// (unlike `write_batch`, whose `?` short-circuits on the first failure)
305    /// and return one `Ok`/`Err` per record. Without this
306    /// override the default impl would surface the first error as an outer
307    /// `Err`, and under `on_batch_error: dlq_all` the pipeline would route the
308    /// *entire* batch to the DLQ — duplicating the already-delivered rows
309    /// against a non-idempotent endpoint (#146 M14).
310    ///
311    /// In **Array** mode the page is POSTed chunk-by-chunk (`batch_size`
312    /// slices), so forward progress is *not* atomic across the whole page —
313    /// each chunk is a separate, independently-committed array POST. The
314    /// override is therefore **chunk-aware** rather than all-or-nothing: it
315    /// iterates the chunks itself and POSTs each array; a row whose chunk was
316    /// delivered is reported `Ok(())`, while the rows of the first failing
317    /// chunk (and every not-yet-sent chunk after it) are reported `Err`.
318    ///
319    /// This is the fix for the duplicate-data bug (F15 / audit #264): the old
320    /// implementation delegated to the all-or-nothing `write_batch`, so a late
321    /// chunk failure surfaced the *whole* page as an outer `Err`. Under
322    /// `on_batch_error: dlq_all` the router then dead-lettered every row —
323    /// including rows from earlier chunks already successfully delivered to the
324    /// live endpoint — producing silent downstream duplicates. By reporting
325    /// per-row outcomes, an already-delivered row is **never** marked failed and
326    /// so can never land in the DLQ.
327    ///
328    /// Within a single failed chunk a single array POST cannot attribute the
329    /// failure to specific rows, so all rows of *that* chunk are reported `Err`
330    /// (acceptable — none of them were delivered). When the **first** chunk
331    /// fails (nothing has been delivered yet) the override preserves the
332    /// original all-or-nothing contract and surfaces an outer `Err`, so the
333    /// router's `on_batch_error` policy (abort vs. dead-letter) still applies to
334    /// a wholly-undelivered page exactly as before.
335    async fn write_batch_partial(
336        &self,
337        records: &[Value],
338    ) -> Result<Vec<faucet_core::RowOutcome>, FaucetError> {
339        if records.is_empty() {
340            return Ok(Vec::new());
341        }
342
343        let auth = self.resolve_auth().await?;
344
345        match &self.config.batch_mode {
346            HttpBatchMode::Individual => {
347                let concurrency = self.config.concurrency.max(1);
348                let auth = &auth;
349                // Attempt every record (failures don't short-circuit the
350                // siblings) with at most `concurrency` POSTs in flight. Tag each
351                // outcome with its index so we can restore record order after
352                // the unordered completion. The per-record futures are built
353                // eagerly (lazy, not yet polled) so `buffer_unordered` drives a
354                // single concrete future type.
355                let pending: Vec<_> =
356                    records
357                        .iter()
358                        .enumerate()
359                        .map(|(idx, record)| async move {
360                            (idx, self.send_with_retry(record, auth).await)
361                        })
362                        .collect();
363                let mut indexed: Vec<(usize, faucet_core::RowOutcome)> =
364                    futures::stream::iter(pending)
365                        .buffer_unordered(concurrency)
366                        .collect()
367                        .await;
368                indexed.sort_by_key(|(idx, _)| *idx);
369                tracing::debug!(
370                    records = records.len(),
371                    "HTTP individual partial batch written"
372                );
373                Ok(indexed.into_iter().map(|(_, outcome)| outcome).collect())
374            }
375            HttpBatchMode::Array => {
376                // `batch_size = 0` is the "no batching" sentinel: forward the
377                // whole page as a single array POST (one chunk). Otherwise
378                // re-chunk into `batch_size` slices and POST one array per
379                // chunk — mirroring `write_batch`, but tracking per-chunk
380                // delivery so a late failure doesn't poison earlier chunks that
381                // were already delivered.
382                let effective_chunk = if self.config.batch_size == 0 {
383                    records.len()
384                } else {
385                    self.config.batch_size
386                };
387
388                let mut outcomes: Vec<faucet_core::RowOutcome> = Vec::with_capacity(records.len());
389                let mut delivered = 0usize;
390                let mut chunks = records.chunks(effective_chunk);
391                let mut failed_chunk: Option<FaucetError> = None;
392
393                for chunk in chunks.by_ref() {
394                    let array = Value::Array(chunk.to_vec());
395                    match self.send_with_retry(&array, &auth).await {
396                        Ok(()) => {
397                            // This chunk was delivered to the live endpoint.
398                            outcomes.extend(chunk.iter().map(|_| Ok(())));
399                            delivered += chunk.len();
400                        }
401                        Err(e) => {
402                            // First failing chunk before any delivery: preserve
403                            // the original all-or-nothing contract so the
404                            // router's `on_batch_error` policy still governs a
405                            // wholly-undelivered page.
406                            if delivered == 0 {
407                                return Err(e);
408                            }
409                            // Otherwise some earlier chunk(s) were delivered;
410                            // mark this chunk's rows (and all remaining,
411                            // never-sent chunks) failed without poisoning the
412                            // delivered rows.
413                            failed_chunk = Some(e);
414                            outcomes.extend(chunk.iter().map(|_| {
415                                Err(FaucetError::Sink(
416                                    "array-mode chunk POST failed; rows not delivered".into(),
417                                ))
418                            }));
419                            break;
420                        }
421                    }
422                }
423
424                if let Some(e) = failed_chunk {
425                    // Remaining chunks were never sent — report them failed too
426                    // so the DLQ captures every undelivered row.
427                    let msg = e.to_string();
428                    for chunk in chunks {
429                        outcomes.extend(chunk.iter().map(|_| {
430                            Err(FaucetError::Sink(format!(
431                                "array-mode chunk not sent after earlier failure: {msg}"
432                            )))
433                        }));
434                    }
435                }
436
437                debug_assert_eq!(
438                    outcomes.len(),
439                    records.len(),
440                    "one outcome per record in array mode"
441                );
442                tracing::debug!(
443                    delivered,
444                    records = records.len(),
445                    batch_size = self.config.batch_size,
446                    "HTTP array partial batch written"
447                );
448                Ok(outcomes)
449            }
450        }
451    }
452}
453
454#[cfg(test)]
455mod tests {
456    use super::*;
457    use crate::config::HttpSinkConfig;
458    use faucet_core::Sink as _;
459
460    #[test]
461    fn dataset_uri_redacts_credentials() {
462        let config = HttpSinkConfig::new("https://user:secret@api.example.com/ingest");
463        let sink = HttpSink::new(config);
464        assert_eq!(sink.dataset_uri(), "https://api.example.com/ingest");
465    }
466
467    #[test]
468    fn creates_sink() {
469        let config = HttpSinkConfig::new("https://api.example.com/ingest");
470        let _sink = HttpSink::new(config);
471    }
472
473    #[test]
474    fn http_sink_is_not_idempotent() {
475        // F32: in Array mode `write_batch` POSTs chunk-by-chunk (non-atomic
476        // forward progress). The HTTP sink must report it does NOT support
477        // idempotent writes, so the pipeline's retry gate (F29) never replays a
478        // partially-delivered page — which would re-POST already-delivered
479        // chunks and silently duplicate rows against the live endpoint.
480        let array = HttpSink::new(
481            HttpSinkConfig::new("https://api.example.com/ingest")
482                .batch_mode(crate::config::HttpBatchMode::Array),
483        );
484        assert!(!array.supports_idempotent_writes());
485        let individual = HttpSink::new(HttpSinkConfig::new("https://api.example.com/ingest"));
486        assert!(!individual.supports_idempotent_writes());
487    }
488
489    #[test]
490    fn build_request_applies_bearer_auth() {
491        let auth = HttpSinkAuth::Bearer {
492            token: "my-token".into(),
493        };
494        let config = HttpSinkConfig::new("https://api.example.com/ingest").auth(auth.clone());
495        let sink = HttpSink::new(config);
496
497        let req = sink
498            .build_request_with_auth(&serde_json::json!({"test": true}), &auth)
499            .unwrap()
500            .build()
501            .unwrap();
502
503        let auth_header = req
504            .headers()
505            .get("authorization")
506            .unwrap()
507            .to_str()
508            .unwrap();
509        assert!(auth_header.starts_with("Bearer "));
510        assert!(auth_header.contains("my-token"));
511    }
512
513    #[test]
514    fn build_request_applies_basic_auth() {
515        let auth = HttpSinkAuth::Basic {
516            username: "user".into(),
517            password: "pass".into(),
518        };
519        let config = HttpSinkConfig::new("https://api.example.com/ingest").auth(auth.clone());
520        let sink = HttpSink::new(config);
521
522        let req = sink
523            .build_request_with_auth(&serde_json::json!({"test": true}), &auth)
524            .unwrap()
525            .build()
526            .unwrap();
527
528        let auth_header = req
529            .headers()
530            .get("authorization")
531            .unwrap()
532            .to_str()
533            .unwrap();
534        assert!(auth_header.starts_with("Basic "));
535    }
536
537    #[test]
538    fn build_request_uses_configured_method() {
539        let config =
540            HttpSinkConfig::new("https://api.example.com/ingest").method(reqwest::Method::PUT);
541        let sink = HttpSink::new(config);
542
543        let req = sink
544            .build_request_with_auth(&serde_json::json!({"test": true}), &HttpSinkAuth::None)
545            .unwrap()
546            .build()
547            .unwrap();
548
549        assert_eq!(req.method(), reqwest::Method::PUT);
550    }
551}