Skip to main content

vector/sinks/
papertrail.rs

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/// Configuration for the `papertrail` sink.
19#[configurable_component(sink("papertrail", "Deliver log events to Papertrail from SolarWinds."))]
20#[derive(Clone, Debug)]
21#[serde(deny_unknown_fields)]
22pub struct PapertrailConfig {
23    /// The TCP endpoint to send logs to.
24    #[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    /// Configures the send buffer size using the `SO_SNDBUF` option on the socket.
37    send_buffer_bytes: Option<usize>,
38
39    /// The value to use as the `process` in Papertrail.
40    #[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}