Skip to main content

risingwave_connector/sink/
http.rs

1// Copyright 2026 RisingWave Labs
2//
3// Licensed under the Apache License, Version 2.0 (the "License");
4// you may not use this file except in compliance with the License.
5// You may obtain a copy of the License at
6//
7//     http://www.apache.org/licenses/LICENSE-2.0
8//
9// Unless required by applicable law or agreed to in writing, software
10// distributed under the License is distributed on an "AS IS" BASIS,
11// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
12// See the License for the specific language governing permissions and
13// limitations under the License.
14
15use std::collections::BTreeMap;
16
17use anyhow::{Context, anyhow};
18use reqwest::header::{CONTENT_TYPE, HeaderMap, HeaderName, HeaderValue};
19use risingwave_common::array::{Op, StreamChunk};
20use risingwave_common::catalog::Schema;
21use risingwave_common::row::Row;
22use risingwave_common::types::{DataType, JsonbRef, ScalarRefImpl};
23use serde::Deserialize;
24use thiserror_ext::AsReport;
25use with_options::WithOptions;
26
27use crate::enforce_secret::EnforceSecret;
28use crate::sink::log_store::DeliveryFutureManagerAddFuture;
29use crate::sink::writer::{
30    AsyncTruncateLogSinkerOf, AsyncTruncateSinkWriter, AsyncTruncateSinkWriterExt,
31};
32use crate::sink::{Result, SINK_TYPE_APPEND_ONLY, Sink, SinkError, SinkParam, SinkWriterParam};
33
34pub const HTTP_SINK: &str = "http";
35const HTTP_SINK_PAYLOAD_COLUMN: &str = "payload";
36const HTTP_SINK_URL_COLUMN: &str = "url";
37
38#[derive(Clone, Debug, Deserialize, WithOptions)]
39pub struct HttpConfig {
40    /// The endpoint URL to send data to.
41    pub url: Option<String>,
42
43    /// HTTP request method. Supported values are `POST` and `PUT`. Defaults to `POST`.
44    pub method: Option<String>,
45
46    /// Content-Type header value. Defaults to `text/plain` for `varchar` and `application/json`
47    /// for `jsonb`.
48    pub content_type: Option<String>,
49
50    /// Sink type, must be "append-only".
51    pub r#type: String,
52
53    #[serde(flatten)]
54    pub unknown_fields: std::collections::HashMap<String, String>,
55}
56
57crate::impl_sink_unknown_fields!(HttpConfig);
58
59impl EnforceSecret for HttpConfig {}
60
61impl HttpConfig {
62    pub fn from_btreemap(
63        values: BTreeMap<String, String>,
64    ) -> Result<(Self, BTreeMap<String, String>)> {
65        // Extract header.* keys before serde parsing
66        let mut headers = BTreeMap::new();
67        let mut rest = BTreeMap::new();
68        for (k, v) in &values {
69            if let Some(header_name) = k.strip_prefix("header.") {
70                headers.insert(header_name.to_owned(), v.clone());
71            } else {
72                rest.insert(k.clone(), v.clone());
73            }
74        }
75
76        let config = serde_json::from_value::<HttpConfig>(serde_json::to_value(rest).unwrap())
77            .map_err(|e| SinkError::Config(anyhow!(e)))?;
78
79        if config.r#type != SINK_TYPE_APPEND_ONLY {
80            return Err(SinkError::Config(anyhow!(
81                "HTTP sink only supports append-only mode"
82            )));
83        }
84
85        Ok((config, headers))
86    }
87}
88
89#[derive(Clone, Debug)]
90enum HttpUrl {
91    Static(reqwest::Url),
92    Dynamic { url_index: usize },
93}
94
95/// Validates the HTTP sink parameters and returns the sink so callers can use it directly without
96/// re-parsing.
97fn validate_http_sink(
98    is_append_only: bool,
99    ignore_delete: bool,
100    schema: &Schema,
101    url: Option<&str>,
102    method: Option<&str>,
103    content_type: Option<&str>,
104    headers: &BTreeMap<String, String>,
105    unknown_fields: std::collections::HashMap<String, String>,
106) -> Result<HttpSink> {
107    if !is_append_only && !ignore_delete {
108        return Err(SinkError::Config(anyhow!(
109            "HTTP sink only supports append-only mode"
110        )));
111    }
112
113    let method = match method {
114        None => reqwest::Method::POST,
115        Some(method) if method.eq_ignore_ascii_case("POST") => reqwest::Method::POST,
116        Some(method) if method.eq_ignore_ascii_case("PUT") => reqwest::Method::PUT,
117        Some(method) => {
118            return Err(SinkError::Config(anyhow!(
119                "HTTP sink method must be POST or PUT, got '{method}'"
120            )));
121        }
122    };
123
124    let fields = schema.fields();
125    let (payload_index, url, payload_type) = if fields.len() == 1 {
126        let Some(url) = url else {
127            return Err(SinkError::Config(anyhow!(
128                "HTTP sink requires url option when schema has exactly 1 column"
129            )));
130        };
131        let url = url
132            .parse()
133            .context("invalid URL")
134            .map_err(SinkError::Config)?;
135        (0, HttpUrl::Static(url), fields[0].data_type.clone())
136    } else {
137        for field in fields {
138            match field.name.as_str() {
139                HTTP_SINK_PAYLOAD_COLUMN | HTTP_SINK_URL_COLUMN => {}
140                _ => {
141                    return Err(SinkError::Config(anyhow!(
142                        "HTTP sink with multiple columns only supports payload and url columns, got {}",
143                        field.name
144                    )));
145                }
146            }
147        }
148
149        let payload_index = fields
150            .iter()
151            .position(|field| field.name == HTTP_SINK_PAYLOAD_COLUMN)
152            .ok_or_else(|| {
153                SinkError::Config(anyhow!(
154                    "HTTP sink with multiple columns requires a payload column"
155                ))
156            })?;
157        let url_index = fields
158            .iter()
159            .position(|field| field.name == HTTP_SINK_URL_COLUMN);
160        let url = match (url, url_index) {
161            (Some(_), Some(_)) => {
162                return Err(SinkError::Config(anyhow!(
163                    "HTTP sink url option cannot coexist with url column"
164                )));
165            }
166            (Some(url), None) => {
167                let url = url
168                    .parse()
169                    .context("invalid URL")
170                    .map_err(SinkError::Config)?;
171                HttpUrl::Static(url)
172            }
173            (None, Some(url_index)) => {
174                if fields[url_index].data_type != DataType::Varchar {
175                    return Err(SinkError::Config(anyhow!(
176                        "HTTP sink url column must be varchar, got {:?}",
177                        fields[url_index].data_type
178                    )));
179                }
180                HttpUrl::Dynamic { url_index }
181            }
182            (None, None) => {
183                return Err(SinkError::Config(anyhow!(
184                    "HTTP sink requires either url option or url column"
185                )));
186            }
187        };
188
189        (payload_index, url, fields[payload_index].data_type.clone())
190    };
191
192    if payload_type != DataType::Varchar && payload_type != DataType::Jsonb {
193        return Err(SinkError::Config(anyhow!(
194            "HTTP sink payload column must be varchar or jsonb, got {:?}",
195            payload_type
196        )));
197    }
198
199    let mut header_map = HeaderMap::new();
200    header_map.insert(
201        CONTENT_TYPE,
202        content_type
203            .unwrap_or(match payload_type {
204                DataType::Varchar => "text/plain",
205                DataType::Jsonb => "application/json",
206                _ => unreachable!("validated HTTP sink column type"),
207            })
208            .parse()
209            .context("invalid content_type")
210            .map_err(SinkError::Config)?,
211    );
212    for (k, v) in headers {
213        let name: HeaderName = k
214            .parse()
215            .with_context(|| format!("invalid header name '{k}'"))
216            .map_err(SinkError::Config)?;
217        let value: HeaderValue = v
218            .parse()
219            .with_context(|| format!("invalid header value for '{k}'"))
220            .map_err(SinkError::Config)?;
221        header_map.insert(name, value);
222    }
223
224    Ok(HttpSink {
225        url,
226        method,
227        payload_index,
228        header_map,
229        unknown_fields,
230    })
231}
232
233#[derive(Clone, Debug)]
234pub struct HttpSink {
235    url: HttpUrl,
236    method: reqwest::Method,
237    payload_index: usize,
238    header_map: HeaderMap,
239    unknown_fields: std::collections::HashMap<String, String>,
240}
241
242impl EnforceSecret for HttpSink {
243    fn enforce_secret<'a>(
244        prop_iter: impl Iterator<Item = &'a str>,
245    ) -> crate::error::ConnectorResult<()> {
246        for prop in prop_iter {
247            HttpConfig::enforce_one(prop)?;
248        }
249        Ok(())
250    }
251}
252
253impl TryFrom<SinkParam> for HttpSink {
254    type Error = SinkError;
255
256    fn try_from(param: SinkParam) -> std::result::Result<Self, Self::Error> {
257        let schema = param.schema();
258        let (config, headers) = HttpConfig::from_btreemap(param.properties)?;
259        validate_http_sink(
260            param.sink_type.is_append_only(),
261            param.ignore_delete,
262            &schema,
263            config.url.as_deref(),
264            config.method.as_deref(),
265            config.content_type.as_deref(),
266            &headers,
267            config.unknown_fields,
268        )
269    }
270}
271
272impl Sink for HttpSink {
273    type LogSinker = AsyncTruncateLogSinkerOf<HttpSinkWriter>;
274
275    const SINK_NAME: &'static str = HTTP_SINK;
276
277    fn validate_unknown_fields(&self) -> Result<()> {
278        crate::sink::validate_sink_unknown_fields(&self.unknown_fields)
279    }
280
281    async fn validate(&self) -> Result<()> {
282        Ok(())
283    }
284
285    async fn new_log_sinker(&self, _writer_param: SinkWriterParam) -> Result<Self::LogSinker> {
286        Ok(HttpSinkWriter::new(
287            self.url.clone(),
288            self.method.clone(),
289            self.payload_index,
290            self.header_map.clone(),
291        )?
292        .into_log_sinker(usize::MAX))
293    }
294}
295
296pub struct HttpSinkWriter {
297    client: reqwest::Client,
298    url: HttpUrl,
299    method: reqwest::Method,
300    payload_index: usize,
301}
302
303impl HttpSinkWriter {
304    fn new(
305        url: HttpUrl,
306        method: reqwest::Method,
307        payload_index: usize,
308        header_map: HeaderMap,
309    ) -> Result<Self> {
310        let client = reqwest::Client::builder()
311            .default_headers(header_map)
312            .build()
313            .context("failed to build HTTP client")
314            .map_err(SinkError::Http)?;
315
316        Ok(Self {
317            client,
318            url,
319            method,
320            payload_index,
321        })
322    }
323
324    fn extract_url(&self, row: &impl Row) -> Result<Option<reqwest::Url>> {
325        match &self.url {
326            HttpUrl::Static(url) => Ok(Some(url.clone())),
327            HttpUrl::Dynamic { url_index } => match row.datum_at(*url_index) {
328                Some(ScalarRefImpl::Utf8(url)) if !url.is_empty() => {
329                    match url.parse::<reqwest::Url>() {
330                        Ok(url) => Ok(Some(url)),
331                        Err(err) => {
332                            tracing::warn!(
333                                error = %err.as_report(),
334                                payload = %self.strip_payload_for_log(row),
335                                "skip HTTP sink row due to invalid URL in url column"
336                            );
337                            Ok(None)
338                        }
339                    }
340                }
341                Some(ScalarRefImpl::Utf8(_)) | None => {
342                    tracing::warn!(
343                        payload = %self.strip_payload_for_log(row),
344                        "skip HTTP sink row due to null or empty url column"
345                    );
346                    Ok(None)
347                }
348                Some(_) => Err(SinkError::Http(anyhow!(
349                    "unexpected url column type, expected varchar"
350                ))),
351            },
352        }
353    }
354
355    fn strip_payload_for_log(&self, row: &impl Row) -> String {
356        match row.datum_at(self.payload_index) {
357            Some(ScalarRefImpl::Utf8(s)) => strip_text_payload(s),
358            Some(ScalarRefImpl::Jsonb(j)) => strip_jsonb_payload(j),
359            Some(_) => "<unexpected payload type>".to_owned(),
360            None => "NULL".to_owned(),
361        }
362    }
363
364    fn extract_payload(&self, row: &impl Row) -> Result<Option<String>> {
365        Ok(match row.datum_at(self.payload_index) {
366            Some(ScalarRefImpl::Utf8(s)) => Some(s.to_owned()),
367            Some(ScalarRefImpl::Jsonb(j)) => Some(j.to_string()),
368            Some(_) => {
369                return Err(SinkError::Http(anyhow!(
370                    "unexpected payload column type, expected varchar or jsonb"
371                )));
372            }
373            None => None, // skip NULL rows
374        })
375    }
376}
377
378fn strip_text_payload(payload: &str) -> String {
379    const EDGE_CHAR_COUNT: usize = 100;
380    let char_count = payload.chars().count();
381    if char_count <= EDGE_CHAR_COUNT * 2 {
382        return payload.to_owned();
383    }
384
385    let prefix: String = payload.chars().take(EDGE_CHAR_COUNT).collect();
386    let suffix: String = payload.chars().skip(char_count - EDGE_CHAR_COUNT).collect();
387    format!("{prefix}...{suffix}")
388}
389
390fn strip_jsonb_payload(payload: JsonbRef<'_>) -> String {
391    if payload.is_array() {
392        let len = payload.array_len().expect("checked JSON array type");
393        return match len {
394            0 => "[]".to_owned(),
395            1 => {
396                let first = payload.access_array_element(0).expect("JSON array element");
397                format!("[{}]", strip_jsonb_payload(first))
398            }
399            2 => {
400                let first = payload.access_array_element(0).expect("JSON array element");
401                let second = payload.access_array_element(1).expect("JSON array element");
402                format!(
403                    "[{},{}]",
404                    strip_jsonb_payload(first),
405                    strip_jsonb_payload(second)
406                )
407            }
408            _ => {
409                let first = payload.access_array_element(0).expect("JSON array element");
410                let last = payload
411                    .access_array_element(len - 1)
412                    .expect("JSON array element");
413                format!(
414                    "[{},...,{}]",
415                    strip_jsonb_payload(first),
416                    strip_jsonb_payload(last)
417                )
418            }
419        };
420    }
421
422    if let Ok(fields) = payload.object_key_values() {
423        let fields = fields
424            .map(|(key, value)| {
425                let key = serde_json::to_string(key).expect("serialize JSON object key");
426                format!("{key}:{}", strip_jsonb_payload(value))
427            })
428            .collect::<Vec<_>>();
429        return format!("{{{}}}", fields.join(","));
430    }
431
432    if let Ok(value) = payload.as_str() {
433        return serde_json::to_string(&strip_text_payload(value)).expect("serialize JSON string");
434    }
435
436    payload.to_string()
437}
438
439impl AsyncTruncateSinkWriter for HttpSinkWriter {
440    async fn write_chunk<'a>(
441        &'a mut self,
442        chunk: StreamChunk,
443        _add_future: DeliveryFutureManagerAddFuture<'a, Self::DeliveryFuture>,
444    ) -> Result<()> {
445        for (op, row) in chunk.rows() {
446            if op != Op::Insert {
447                continue;
448            }
449
450            let Some(payload) = self.extract_payload(&row)? else {
451                continue;
452            };
453            let Some(url) = self.extract_url(&row)? else {
454                continue;
455            };
456
457            let resp = self
458                .client
459                .request(self.method.clone(), url)
460                .body(payload)
461                .send()
462                .await
463                .context("HTTP request failed")
464                .map_err(SinkError::Http)?;
465
466            if !resp.status().is_success() {
467                let status = resp.status();
468                let body = resp.text().await.unwrap_or_default();
469                return Err(SinkError::Http(anyhow!(
470                    "HTTP sink received non-success response: {} {}",
471                    status,
472                    body
473                )));
474            }
475        }
476
477        Ok(())
478    }
479}
480
481#[cfg(test)]
482mod tests {
483    use risingwave_common::types::{JsonbVal, Scalar};
484
485    use super::*;
486
487    #[test]
488    fn test_strip_text_payload() {
489        assert_eq!(strip_text_payload("short payload"), "short payload");
490
491        let payload = format!("{}{}", "a".repeat(100), "z".repeat(100));
492        assert_eq!(strip_text_payload(&payload), payload);
493
494        let payload = format!("{}middle{}", "a".repeat(101), "z".repeat(101));
495        assert_eq!(
496            strip_text_payload(&payload),
497            format!("{}...{}", "a".repeat(100), "z".repeat(100))
498        );
499    }
500
501    #[test]
502    fn test_strip_jsonb_payload() {
503        let payload: JsonbVal = r#"{
504            "id": 1,
505            "items": [1, {"nested": ["first", "middle", "last"]}, 3],
506            "message": "short"
507        }"#
508        .parse()
509        .unwrap();
510
511        assert_eq!(
512            strip_jsonb_payload(payload.as_scalar_ref()),
513            r#"{"id":1,"items":[1,...,3],"message":"short"}"#
514        );
515
516        let payload: JsonbVal = r#"["only"]"#.parse().unwrap();
517        assert_eq!(strip_jsonb_payload(payload.as_scalar_ref()), r#"["only"]"#);
518
519        let payload: JsonbVal = r#"["first","last"]"#.parse().unwrap();
520        assert_eq!(
521            strip_jsonb_payload(payload.as_scalar_ref()),
522            r#"["first","last"]"#
523        );
524    }
525}