use serde_json::json;
use wp_connector_api::{ConnectorDef, ConnectorScope, ParamMap};
pub fn builtin_sink_defs() -> Vec<ConnectorDef> {
let mut defs = Vec::new();
{
let mut params = ParamMap::new();
params.insert("sleep_ms".into(), json!(0));
defs.push(ConnectorDef {
id: "blackhole_sink".into(),
kind: "blackhole".into(),
scope: ConnectorScope::Sink,
allow_override: vec!["sleep_ms".into(), "protocol".into()],
default_params: params,
origin: Some("builtin:blackhole".into()),
});
}
{
let mut params = ParamMap::new();
params.insert("fmt".into(), json!("json"));
params.insert("base".into(), json!("./data/out_dat"));
params.insert("file".into(), json!("default.json"));
params.insert("sync".into(), json!(false));
defs.push(ConnectorDef {
id: "file_json_sink".into(),
kind: "file".into(),
scope: ConnectorScope::Sink,
allow_override: vec![
"base".into(),
"file".into(),
"sync".into(),
"protocol".into(),
],
default_params: params,
origin: Some("builtin:file".into()),
});
}
{
let mut params = ParamMap::new();
params.insert("fmt".into(), json!("proto-text"));
params.insert("base".into(), json!("./data/out_dat"));
params.insert("file".into(), json!("default.pbtxt"));
params.insert("sync".into(), json!(false));
defs.push(ConnectorDef {
id: "file_proto_text_sink".into(),
kind: "file".into(),
scope: ConnectorScope::Sink,
allow_override: vec![
"base".into(),
"file".into(),
"sync".into(),
"protocol".into(),
],
default_params: params,
origin: Some("builtin:file".into()),
});
}
{
let mut params = ParamMap::new();
params.insert("fmt".into(), json!("proto-text"));
params.insert("base".into(), json!("./data/out_dat"));
params.insert("file".into(), json!("default.dat"));
params.insert("sync".into(), json!(false));
defs.push(ConnectorDef {
id: "file_proto_sink".into(),
kind: "file".into(),
scope: ConnectorScope::Sink,
allow_override: vec![
"base".into(),
"file".into(),
"sync".into(),
"protocol".into(),
],
default_params: params,
origin: Some("builtin:file".into()),
});
}
{
let mut params = ParamMap::new();
params.insert("fmt".into(), json!("kv"));
params.insert("base".into(), json!("./data/out_dat"));
params.insert("file".into(), json!("default.kv"));
params.insert("sync".into(), json!(false));
defs.push(ConnectorDef {
id: "file_kv_sink".into(),
kind: "file".into(),
scope: ConnectorScope::Sink,
allow_override: vec![
"base".into(),
"file".into(),
"sync".into(),
"protocol".into(),
],
default_params: params,
origin: Some("builtin:file".into()),
});
}
{
let mut params = ParamMap::new();
params.insert("fmt".into(), json!("raw"));
params.insert("base".into(), json!("./data/out_dat"));
params.insert("file".into(), json!("default.raw"));
params.insert("sync".into(), json!(false));
defs.push(ConnectorDef {
id: "file_raw_sink".into(),
kind: "file".into(),
scope: ConnectorScope::Sink,
allow_override: vec![
"base".into(),
"file".into(),
"sync".into(),
"protocol".into(),
],
default_params: params,
origin: Some("builtin:file".into()),
});
}
{
let mut params = ParamMap::new();
params.insert("addr".into(), json!("127.0.0.1"));
params.insert("port".into(), json!(1514));
params.insert("protocol".into(), json!("udp"));
params.insert("strip_header".into(), json!(true));
params.insert("attach_meta_tags".into(), json!(true));
params.insert("tcp_recv_bytes".into(), json!(256000));
defs.push(ConnectorDef {
id: "syslog_sink".into(),
kind: "syslog".into(),
scope: ConnectorScope::Sink,
allow_override: vec![
"addr".into(),
"port".into(),
"protocol".into(),
"app_name".into(),
],
default_params: params,
origin: Some("builtin:syslog_sink".into()),
});
}
{
let mut params = ParamMap::new();
params.insert("addr".into(), json!("127.0.0.1"));
params.insert("port".into(), json!(9000));
params.insert("framing".into(), json!("line"));
defs.push(ConnectorDef {
id: "tcp_sink".into(),
kind: "tcp".into(),
scope: ConnectorScope::Sink,
allow_override: vec![
"addr".into(),
"port".into(),
"framing".into(),
"protocol".into(),
],
default_params: params,
origin: Some("builtin:tcp_sink".into()),
});
}
{
let mut params = ParamMap::new();
params.insert("protocol".into(), json!("arrow"));
params.insert("base".into(), json!("./data/out_dat"));
params.insert("file".into(), json!("default.arrow"));
params.insert("sync".into(), json!(false));
defs.push(ConnectorDef {
id: "file_arrow_sink".into(),
kind: "file".into(),
scope: ConnectorScope::Sink,
allow_override: vec![
"base".into(),
"file".into(),
"sync".into(),
"protocol".into(),
],
default_params: params,
origin: Some("builtin:file_arrow_sink".into()),
});
}
{
let mut params = ParamMap::new();
params.insert("protocol".into(), json!("arrow"));
params.insert("addr".into(), json!("127.0.0.1"));
params.insert("port".into(), json!(9000));
defs.push(ConnectorDef {
id: "tcp_arrow_sink".into(),
kind: "tcp".into(),
scope: ConnectorScope::Sink,
allow_override: vec!["addr".into(), "port".into(), "protocol".into()],
default_params: params,
origin: Some("builtin:tcp_arrow_sink".into()),
});
}
{
let mut params = ParamMap::new();
params.insert("fmt".into(), json!("kv"));
params.insert("base".into(), json!("./data/out_dat"));
params.insert("file".into(), json!("default.kv"));
defs.push(ConnectorDef {
id: "file_rescue_sink".into(),
kind: "test_rescue".into(),
scope: ConnectorScope::Sink,
allow_override: vec!["base".into(), "file".into(), "protocol".into()],
default_params: params,
origin: Some("builtin:test_rescue".into()),
});
}
defs
}
pub fn builtin_source_defs() -> Vec<ConnectorDef> {
let mut defs = Vec::new();
{
let mut params = ParamMap::new();
params.insert("base".into(), json!("./data/in_dat"));
params.insert("file".into(), json!("gen.dat"));
params.insert("encode".into(), json!("text"));
params.insert("data_format".into(), json!("ndjson"));
defs.push(ConnectorDef {
id: "file_src".into(),
kind: "file".into(),
scope: ConnectorScope::Source,
allow_override: vec![
"base".into(),
"file".into(),
"encode".into(),
"data_format".into(),
],
default_params: params,
origin: Some("builtin:file_source".into()),
});
}
{
let mut params = ParamMap::new();
params.insert("addr".into(), json!("0.0.0.0"));
params.insert("port".into(), json!(514));
params.insert("protocol".into(), json!("udp"));
params.insert("tcp_recv_bytes".into(), json!(10_485_760));
params.insert("udp_recv_buffer".into(), json!(8_388_608));
params.insert("header_mode".into(), json!("skip"));
params.insert("fast_strip".into(), json!(false));
defs.push(ConnectorDef {
id: "syslog_src".into(),
kind: "syslog".into(),
scope: ConnectorScope::Source,
allow_override: vec![
"addr".into(),
"port".into(),
"protocol".into(),
"tcp_recv_bytes".into(),
"udp_recv_buffer".into(),
"header_mode".into(),
"fast_strip".into(),
],
default_params: params,
origin: Some("builtin:syslog_source".into()),
});
}
{
let mut params = ParamMap::new();
params.insert("addr".into(), json!("0.0.0.0"));
params.insert("port".into(), json!(9000));
params.insert("framing".into(), json!("auto"));
params.insert("tcp_recv_bytes".into(), json!(256_000));
params.insert("instances".into(), json!(1));
params.insert("data_format".into(), json!("ndjson"));
defs.push(ConnectorDef {
id: "tcp_src".into(),
kind: "tcp".into(),
scope: ConnectorScope::Source,
allow_override: vec![
"addr".into(),
"port".into(),
"framing".into(),
"tcp_recv_bytes".into(),
"instances".into(),
"data_format".into(),
],
default_params: params,
origin: Some("builtin:tcp_source".into()),
});
}
defs
}
pub fn sink_def(id: &str) -> Option<ConnectorDef> {
builtin_sink_defs().into_iter().find(|d| d.id == id)
}
pub fn source_def(id: &str) -> Option<ConnectorDef> {
builtin_source_defs().into_iter().find(|d| d.id == id)
}
pub fn sink_defs_by_kind(kind: &str) -> Vec<ConnectorDef> {
builtin_sink_defs()
.into_iter()
.filter(|d| d.kind.as_str() == kind)
.collect()
}
#[cfg(test)]
mod tests {
use super::*;
#[test]
fn builtin_file_sink_defs_include_raw_variant() {
let raw = sink_def("file_raw_sink").expect("file_raw_sink builtin def");
assert_eq!(raw.kind, "file");
assert_eq!(raw.scope, ConnectorScope::Sink);
assert_eq!(
raw.default_params
.get("fmt")
.and_then(|value| value.as_str()),
Some("raw")
);
assert_eq!(
raw.default_params
.get("file")
.and_then(|value| value.as_str()),
Some("default.raw")
);
assert!(raw.allow_override.iter().any(|key| key == "base"));
assert!(raw.allow_override.iter().any(|key| key == "file"));
assert!(raw.allow_override.iter().any(|key| key == "sync"));
}
#[test]
fn file_factory_def_provider_exposes_raw_variant() {
let ids: Vec<_> = sink_defs_by_kind("file")
.into_iter()
.map(|def| def.id)
.collect();
assert!(ids.iter().any(|id| id == "file_json_sink"));
assert!(ids.iter().any(|id| id == "file_raw_sink"));
}
#[test]
fn all_sink_defs_allow_override_protocol() {
for def in builtin_sink_defs() {
assert!(
def.allow_override.iter().any(|k| k == "protocol"),
"{} should allow override of 'protocol'",
def.id
);
}
}
}