1use bytes::{BufMut, BytesMut};
2use syslog::{Facility, Formatter3164, LogFormat, Severity};
3use vector_lib::configurable::configurable_component;
4use vrl::value::Kind;
5
6use crate::{
7 codecs::{Encoder, EncodingConfig, Transformer},
8 config::{AcknowledgementsConfig, DataType, GenerateConfig, Input, SinkConfig, SinkContext},
9 event::Event,
10 internal_events::TemplateRenderingError,
11 schema,
12 sinks::util::{UriSerde, tcp::TcpSinkConfig},
13 tcp::TcpKeepaliveConfig,
14 template::UnconfinedTemplate,
15 tls::TlsEnableableConfig,
16};
17
18#[configurable_component(sink("papertrail", "Deliver log events to Papertrail from SolarWinds."))]
20#[derive(Clone, Debug)]
21#[serde(deny_unknown_fields)]
22pub struct PapertrailConfig {
23 #[configurable(metadata(docs::examples = "logs.papertrailapp.com:12345"))]
25 endpoint: UriSerde,
26
27 #[configurable(derived)]
28 encoding: EncodingConfig,
29
30 #[configurable(derived)]
31 keepalive: Option<TcpKeepaliveConfig>,
32
33 #[configurable(derived)]
34 tls: Option<TlsEnableableConfig>,
35
36 send_buffer_bytes: Option<usize>,
38
39 #[configurable(metadata(docs::examples = "{{ process }}", docs::examples = "my-process",))]
41 #[serde(default = "default_process")]
42 process: UnconfinedTemplate,
43
44 #[configurable(derived)]
45 #[serde(
46 default,
47 deserialize_with = "crate::serde::bool_or_struct",
48 skip_serializing_if = "crate::serde::is_default"
49 )]
50 acknowledgements: AcknowledgementsConfig,
51}
52
53fn default_process() -> UnconfinedTemplate {
54 UnconfinedTemplate::try_from("vector").unwrap()
55}
56
57impl GenerateConfig for PapertrailConfig {
58 fn generate_config() -> serde_json::Value {
59 serde_yaml::from_str(indoc::indoc! {
60 r#"endpoint: "logs.papertrailapp.com:12345"
61 encoding:
62 codec: json"#,
63 })
64 .unwrap()
65 }
66}
67
68#[async_trait::async_trait]
69#[typetag::serde(name = "papertrail")]
70impl SinkConfig for PapertrailConfig {
71 async fn build(
72 &self,
73 _cx: SinkContext,
74 ) -> crate::Result<(super::VectorSink, super::Healthcheck)> {
75 let host = self
76 .endpoint
77 .uri
78 .host()
79 .map(str::to_string)
80 .ok_or_else(|| "A host is required for endpoint".to_string())?;
81 let port = self
82 .endpoint
83 .uri
84 .port_u16()
85 .ok_or_else(|| "A port is required for endpoint".to_string())?;
86
87 let address = format!("{host}:{port}");
88 let tls = Some(
89 self.tls
90 .clone()
91 .unwrap_or_else(TlsEnableableConfig::enabled),
92 );
93
94 let pid = std::process::id();
95 let process = self.process.clone();
96
97 let sink_config = TcpSinkConfig::new(address, self.keepalive, tls, self.send_buffer_bytes);
98
99 let transformer = self.encoding.transformer();
100 let serializer = self.encoding.build()?;
101 let encoder = Encoder::<()>::new(serializer);
102
103 sink_config.build(
104 Transformer::default(),
105 PapertrailEncoder {
106 pid,
107 process,
108 transformer,
109 encoder,
110 },
111 )
112 }
113
114 fn input(&self) -> Input {
115 let requirement = schema::Requirement::empty().optional_meaning("host", Kind::bytes());
116
117 Input::new(self.encoding.config().input_type() & DataType::Log)
118 .with_schema_requirement(requirement)
119 }
120
121 fn acknowledgements(&self) -> &AcknowledgementsConfig {
122 &self.acknowledgements
123 }
124}
125
126#[derive(Debug, Clone)]
127struct PapertrailEncoder {
128 pid: u32,
129 process: UnconfinedTemplate,
130 transformer: Transformer,
131 encoder: Encoder<()>,
132}
133
134impl tokio_util::codec::Encoder<Event> for PapertrailEncoder {
135 type Error = vector_lib::codecs::encoding::Error;
136
137 fn encode(
138 &mut self,
139 mut event: Event,
140 buffer: &mut bytes::BytesMut,
141 ) -> Result<(), Self::Error> {
142 let host = event
143 .as_mut_log()
144 .get_host()
145 .map(|host| host.to_string_lossy().into_owned());
146
147 let process = self
148 .process
149 .render_string(&event)
150 .map_err(|error| {
151 emit!(TemplateRenderingError {
152 error,
153 field: Some("process"),
154 drop_event: false,
155 })
156 })
157 .ok()
158 .unwrap_or_else(|| String::from("vector"));
159
160 let formatter = Formatter3164 {
161 facility: Facility::LOG_USER,
162 hostname: host,
163 process,
164 pid: self.pid,
165 };
166
167 self.transformer.transform(&mut event);
168
169 let mut bytes = BytesMut::new();
170 self.encoder.encode(event, &mut bytes)?;
171
172 let message = String::from_utf8_lossy(&bytes);
173
174 formatter
175 .format(&mut buffer.writer(), Severity::LOG_INFO, message)
176 .map_err(|error| Self::Error::SerializingError(format!("{error}").into()))?;
177
178 buffer.put_u8(b'\n');
179
180 Ok(())
181 }
182}
183
184#[cfg(test)]
185mod tests {
186 use bytes::BytesMut;
187 use futures::{future::ready, stream};
188 use tokio_util::codec::Encoder as _;
189 use vector_lib::{
190 codecs::JsonSerializerConfig,
191 event::{Event, LogEvent},
192 };
193 use vrl::event_path;
194
195 use super::*;
196 use crate::test_util::{
197 components::{SINK_TAGS, run_and_assert_sink_compliance},
198 http::{always_200_response, spawn_blackhole_http_server},
199 };
200
201 #[test]
202 fn generate_config() {
203 crate::test_util::test_generate_config::<PapertrailConfig>();
204 }
205
206 #[tokio::test]
207 async fn component_spec_compliance() {
208 let mock_endpoint = spawn_blackhole_http_server(always_200_response).await;
209
210 let mut config: PapertrailConfig =
211 serde_json::from_value(PapertrailConfig::generate_config())
212 .expect("config should be valid");
213 config.endpoint = mock_endpoint.try_into().unwrap();
214 config.tls = Some(TlsEnableableConfig::default());
215
216 let context = SinkContext::default();
217 let (sink, _healthcheck) = config.build(context).await.unwrap();
218
219 let event = Event::Log(LogEvent::from("simple message"));
220 run_and_assert_sink_compliance(sink, stream::once(ready(event)), &SINK_TAGS).await;
221 }
222
223 #[test]
224 fn encode_event_apply_rules() {
225 let mut evt = Event::Log(LogEvent::from("vector"));
226 evt.as_mut_log().insert(event_path!("magic"), "key");
227 evt.as_mut_log().insert(event_path!("process"), "foo");
228
229 let mut encoder = PapertrailEncoder {
230 pid: 0,
231 process: UnconfinedTemplate::try_from("{{ process }}").unwrap(),
232 transformer: Transformer::new(None, Some(vec!["magic".into()]), None).unwrap(),
233 encoder: Encoder::<()>::new(JsonSerializerConfig::default().build().into()),
234 };
235
236 let mut bytes = BytesMut::new();
237 encoder.encode(evt, &mut bytes).unwrap();
238 let bytes = bytes.freeze();
239
240 let msg = bytes.slice(String::from_utf8_lossy(&bytes).find(": ").unwrap() + 2..bytes.len());
241 let value: serde_json::Value = serde_json::from_slice(&msg).unwrap();
242 let value = value.as_object().unwrap();
243
244 assert!(!value.contains_key("magic"));
245 assert_eq!(value.get("process").unwrap().as_str(), Some("foo"));
246 }
247}