1use 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 pub url: Option<String>,
42
43 pub method: Option<String>,
45
46 pub content_type: Option<String>,
49
50 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 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
95fn 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, })
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}