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}