@cryptotaxi247 / netdata-1 / commits / 322bde853

chore(otel/journaldexporter): add socket/remote clients (#20121)

wip socket/remote journald clients

Ilya Mashchenko committed Apr 13, 2025 at 18:30 UTC 322bde853ade2c0cbcb44c233867c7722ebde1e6
8 files changed +1021 -37
src/go/otel-collector/exporter/journaldexporter/config.go
+12 -4
@@ -1,14 +1,22 @@
1 +// SPDX-License-Identifier: GPL-3.0-or-later
2 +
3 package journaldexporter
4
5 import (
6 + "time"
7 +
8 "go.opentelemetry.io/collector/component"
9 )
10
11 type Config struct {
8 - URL string `mapstructure:"url"`
9 - ServerKeyFile string `mapstructure:"server_key_file"`
10 - ServerCertificateFile string `mapstructure:"server_certificate_file"`
11 - TrustedCertificateFile string `mapstructure:"trusted_certificate_file"`
12 + URL string `mapstructure:"url"`
13 + Timeout time.Duration `mapstructure:"timeout"`
14 + TLS struct {
15 + SrvCertFile string `mapstructure:"server_certificate_file"`
16 + SrvKeyFile string `mapstructure:"server_key_file"`
17 + TrustedCertFile string `mapstructure:"trusted_certificate_file"`
18 + InsecureSkipVerify bool `mapstructure:"insecure_skip_verify"`
19 + } `mapstructure:"tls"`
20 }
21
22 var _ component.Config = (*Config)(nil)
src/go/otel-collector/exporter/journaldexporter/convert.go
+15 -30
@@ -1,3 +1,5 @@
1 +// SPDX-License-Identifier: GPL-3.0-or-later
2 +
3 package journaldexporter
4
5 import (
@@ -14,20 +16,14 @@ import (
16
17 func (e *journaldExporter) logsToJournaldMessages(ld plog.Logs, buf *bytes.Buffer) {
18 receivedAt := fmt.Sprintf("%d", time.Now().UnixNano()/1000)
17 - rls := ld.ResourceLogs()
18 - for i := 0; i < rls.Len(); i++ {
19 - rl := rls.At(i)
19 +
20 + for _, rl := range ld.ResourceLogs().All() {
21 resource := rl.Resource()
22
22 - scopeLogs := rl.ScopeLogs()
23 - for j := 0; j < scopeLogs.Len(); j++ {
24 - sl := scopeLogs.At(j)
23 + for _, sl := range rl.ScopeLogs().All() {
24 scope := sl.Scope()
25
27 - logRecords := sl.LogRecords()
28 - for k := 0; k < logRecords.Len(); k++ {
29 - lr := logRecords.At(k)
30 -
26 + for _, lr := range sl.LogRecords().All() {
27 writeField(buf, "__REALTIME_TIMESTAMP", receivedAt)
28 writeField(buf, "SYSLOG_IDENTIFIER", e.fields.syslogID)
29 writeField(buf, "_PID", e.fields.pid)
@@ -38,31 +34,23 @@ func (e *journaldExporter) logsToJournaldMessages(ld plog.Logs, buf *bytes.Buffe
34 writeField(buf, "PRIORITY", strconv.Itoa(mapSeverityToJournaldPriority(lr.SeverityNumber())))
35 writeField(buf, "MESSAGE", bodyToString(lr.Body()))
36
41 - resource.Attributes().Range(func(k string, v pcommon.Value) bool {
37 + for k, v := range resource.Attributes().All() {
38 writeField(buf, "OTEL_RESOURCE_ATTR_"+k, v.AsString())
43 - return true
44 - })
45 -
46 - if scope.Name() != "" {
47 - writeField(buf, "OTEL_SCOPE_NAME", scope.Name())
48 - }
49 - if scope.Version() != "" {
50 - writeField(buf, "OTEL_SCOPE_VERSION", scope.Version())
39 }
40
41 + writeField(buf, "OTEL_SCOPE_NAME", scope.Name())
42 + writeField(buf, "OTEL_SCOPE_VERSION", scope.Version())
43 +
44 if lr.Timestamp() != 0 {
45 ts := time.Unix(0, int64(lr.Timestamp()))
46 writeField(buf, "OTEL_TIMESTAMP", strconv.FormatInt(ts.UnixMicro(), 10))
47 }
57 -
58 - if lr.ObservedTimestamp() != 0 && lr.ObservedTimestamp() != lr.Timestamp() {
48 + if lr.ObservedTimestamp() != 0 {
49 ts := time.Unix(0, int64(lr.ObservedTimestamp()))
50 writeField(buf, "OTEL_OBSERVED_TIMESTAMP", strconv.FormatInt(ts.UnixMicro(), 10))
51 }
52
63 - if lr.SeverityText() != "" {
64 - writeField(buf, "OTEL_SEVERITY_LEVEL", lr.SeverityText())
65 - }
53 + writeField(buf, "OTEL_SEVERITY_LEVEL", lr.SeverityText())
54
55 if !lr.TraceID().IsEmpty() {
56 writeField(buf, "OTEL_TRACE_ID", lr.TraceID().String())
@@ -74,14 +62,11 @@ func (e *journaldExporter) logsToJournaldMessages(ld plog.Logs, buf *bytes.Buffe
62 }
63 }
64
77 - if lr.EventName() != "" {
78 - writeField(buf, "OTEL_EVENT_NAME", lr.EventName())
79 - }
65 + writeField(buf, "OTEL_EVENT_NAME", lr.EventName())
66
81 - lr.Attributes().Range(func(k string, v pcommon.Value) bool {
67 + for k, v := range lr.Attributes().All() {
68 writeField(buf, "OTEL_ATTR_"+k, v.AsString())
83 - return true
84 - })
69 + }
70
71 buf.WriteByte('\n') // extra newline
72 }
src/go/otel-collector/exporter/journaldexporter/convert_test.go
+2
@@ -1,3 +1,5 @@
1 +// SPDX-License-Identifier: GPL-3.0-or-later
2 +
3 package journaldexporter
4
5 import (
src/go/otel-collector/exporter/journaldexporter/exporter.go
+25 -3
@@ -1,6 +1,9 @@
1 +// SPDX-License-Identifier: GPL-3.0-or-later
2 +
3 package journaldexporter
4
5 import (
6 + "bytes"
7 "context"
8 "fmt"
9 "os"
@@ -17,7 +20,14 @@ type (
20 log *zap.Logger
21 conf *Config
22 fields commonFields
23 + s sender
24 + buf bytes.Buffer
25 + }
26 + sender interface {
27 + sendMessage(ctx context.Context, msg []byte) error
28 + shutdown(ctx context.Context) error
29 }
30 +
31 commonFields struct {
32 syslogID string
33 pid string
@@ -35,9 +45,21 @@ func newJournaldExporter(cfg component.Config, logger *zap.Logger) *journaldExpo
45 }
46 }
47
38 -func (e *journaldExporter) consumeLogs(_ context.Context, ld plog.Logs) error {
39 - fmt.Println("consume logs")
40 - return nil
48 +func (e *journaldExporter) consumeLogs(ctx context.Context, ld plog.Logs) error {
49 + if e.s == nil {
50 + return nil
51 + }
52 +
53 + e.buf.Reset()
54 +
55 + e.logsToJournaldMessages(ld, &e.buf)
56 +
57 + select {
58 + case <-ctx.Done():
59 + return nil
60 + default:
61 + return e.s.sendMessage(ctx, e.buf.Bytes())
62 + }
63 }
64
65 func (e *journaldExporter) Start(_ context.Context, _ component.Host) error {
src/go/otel-collector/exporter/journaldexporter/factory.go
+2
@@ -1,3 +1,5 @@
1 +// SPDX-License-Identifier: GPL-3.0-or-later
2 +
3 package journaldexporter
4
5 import (
src/go/otel-collector/exporter/journaldexporter/journal_remote.go new
+287
@@ -0,0 +1,287 @@
1 +package journaldexporter
2 +
3 +import (
4 + "context"
5 + "crypto/tls"
6 + "crypto/x509"
7 + "errors"
8 + "fmt"
9 + "io"
10 + "net"
11 + "net/http"
12 + "os"
13 + "sync"
14 +
15 + "go.uber.org/zap"
16 +)
17 +
18 +func newRemoteJournalClient(cfg *Config) (*remoteJournalClient, error) {
19 + client, err := newHTTPClient(cfg)
20 + if err != nil {
21 + return nil, fmt.Errorf("failed to create remote journal http client: %v", err)
22 + }
23 + return &remoteJournalClient{
24 + remoteURL: cfg.URL,
25 + httpClient: client,
26 + done: make(chan struct{}),
27 + }, nil
28 +}
29 +
30 +type remoteJournalClient struct {
31 + log *zap.Logger
32 +
33 + remoteURL string
34 + httpClient *http.Client
35 +
36 + mu sync.Mutex // Protects access to writer, error info, cancel func
37 + w io.WriteCloser // The pipe writer for the current active upload stream
38 + reqErr error // Stores the error from the last background upload attempt
39 + reqCancel context.CancelFunc // Function to cancel the current background HTTP request
40 +
41 + done chan struct{} // Closed when the httpClient is shutting down
42 + wg sync.WaitGroup // Waits for background operations to complete during done
43 +}
44 +
45 +func (jc *remoteJournalClient) sendMessage(ctx context.Context, msg []byte) error {
46 + if len(msg) == 0 {
47 + return nil
48 + }
49 +
50 + var currentWriter io.WriteCloser
51 + var connectErr error
52 +
53 + jc.mu.Lock()
54 +
55 + select {
56 + case <-ctx.Done():
57 + jc.mu.Unlock()
58 + return ctx.Err()
59 + case <-jc.done:
60 + jc.mu.Unlock()
61 + return errors.New("journal: client is shut down")
62 + default:
63 + }
64 +
65 + if jc.w == nil {
66 + jc.log.Info("journal: not connected, attempting to connect...")
67 + connectErr = jc.connectLocked(ctx)
68 + }
69 +
70 + if connectErr == nil && jc.w != nil {
71 + currentWriter = jc.w
72 + }
73 + lastReqErr := jc.reqErr
74 + jc.mu.Unlock()
75 +
76 + if connectErr != nil {
77 + return fmt.Errorf("journal: connection attempt failed: %w", connectErr)
78 + }
79 +
80 + if currentWriter == nil {
81 + errMsg := "journal: connection unavailable"
82 + if lastReqErr != nil {
83 + return fmt.Errorf("%s (last background error: %w)", errMsg, lastReqErr)
84 + }
85 + return errors.New(errMsg)
86 + }
87 +
88 + if _, err := currentWriter.Write(msg); err != nil {
89 + return fmt.Errorf("journal: failed to write message (connection likely closed): %w", err)
90 + }
91 +
92 + return nil
93 +}
94 +
95 +// connectLocked initiates a new upload stream. Must be called with jc.mu held.
96 +func (jc *remoteJournalClient) connectLocked(ctx context.Context) error {
97 + if jc.w != nil {
98 + jc.log.Warn("journal: warning - connectLocked called while already connected")
99 + jc.disconnectLocked()
100 + }
101 +
102 + if jc.isDone() {
103 + return errors.New("journal: client is shut down")
104 + }
105 +
106 + reqCtx, cancel := context.WithCancel(ctx)
107 + jc.reqCancel = cancel
108 +
109 + pr, pw := io.Pipe()
110 + jc.w = pw
111 + jc.reqErr = nil
112 +
113 + jc.wg.Add(1)
114 + go func() {
115 + defer jc.wg.Done()
116 +
117 + // blocks until the stream finishes.
118 + reqErr := jc.doRequest(reqCtx, pr)
119 +
120 + jc.mu.Lock()
121 + if jc.w == pw {
122 + jc.reqErr = reqErr
123 + jc.w = nil
124 + jc.reqCancel = nil
125 + }
126 + jc.mu.Unlock()
127 +
128 + if reqErr != nil {
129 + if errors.Is(reqErr, context.Canceled) || errors.Is(reqErr, io.ErrClosedPipe) {
130 + jc.log.Info("journal: background upload finished successfully.")
131 + } else {
132 + jc.log.Error("journal: background upload failed with error", zap.Error(reqErr))
133 + }
134 + }
135 + }()
136 +
137 + jc.log.Info("journal: connection attempt initiated (running in background)")
138 +
139 + return nil
140 +}
141 +
142 +func (jc *remoteJournalClient) doRequest(ctx context.Context, pr *io.PipeReader) (err error) {
143 + defer func() { _ = pr.CloseWithError(err) }()
144 +
145 + req, err := http.NewRequestWithContext(ctx, http.MethodPost, jc.remoteURL, pr)
146 + if err != nil {
147 + return fmt.Errorf("could not create request: %w", err)
148 + }
149 +
150 + req.Header.Set("Content-Type", "application/vnd.fdo.journal")
151 +
152 + // This call blocks until the request body (pr) is closed, the context is canceled,
153 + // the server responds AND closes the connection, or a connection error occurs.
154 + resp, err := jc.httpClient.Do(req)
155 + if err != nil {
156 + if errors.Is(err, context.Canceled) {
157 + jc.log.Info("journal: request cancelled via context.")
158 + } else {
159 + jc.log.Warn(fmt.Sprintf("journal: http client error: %v", err))
160 + }
161 + return err
162 + }
163 +
164 + defer closeBody(resp)
165 +
166 + if resp.StatusCode < 200 || resp.StatusCode >= 300 {
167 + bodyBytes, _ := io.ReadAll(resp.Body)
168 + return fmt.Errorf("unexpected status code %d: %s", resp.StatusCode, string(bodyBytes))
169 + }
170 +
171 + return nil
172 +}
173 +
174 +func (jc *remoteJournalClient) shutdown(ctx context.Context) error {
175 + jc.mu.Lock()
176 + if jc.isDone() {
177 + jc.mu.Unlock()
178 + jc.log.Warn("journal: client is shut down")
179 + return nil
180 + }
181 +
182 + close(jc.done)
183 + jc.log.Info("journal: shutdown initiated.")
184 + jc.disconnectLocked()
185 +
186 + jc.mu.Unlock()
187 +
188 + done := make(chan struct{})
189 + go func() {
190 + jc.wg.Wait()
191 + close(done)
192 + }()
193 +
194 + select {
195 + case <-done:
196 + jc.log.Info("journal: all background tasks finished.")
197 + return nil
198 + case <-ctx.Done():
199 + jc.log.Info("journal: shutdown timed out waiting for background tasks.")
200 + return ctx.Err()
201 + }
202 +}
203 +
204 +// disconnectLocked cancels the current request and closes the writer pipe. Must be called with jc.mu held.
205 +func (jc *remoteJournalClient) disconnectLocked() {
206 + if jc.w != nil {
207 + _ = jc.w.Close()
208 + jc.w = nil
209 + }
210 + if jc.reqCancel != nil {
211 + jc.reqCancel()
212 + jc.reqCancel = nil
213 + }
214 +}
215 +
216 +func (jc *remoteJournalClient) isDone() bool {
217 + select {
218 + case <-jc.done:
219 + return true
220 + default:
221 + return false
222 + }
223 +}
224 +
225 +func newHTTPClient(cfg *Config) (*http.Client, error) {
226 + tlsConfig, err := newTLSConfig(cfg)
227 + if err != nil {
228 + return nil, err
229 + }
230 +
231 + d := &net.Dialer{Timeout: cfg.Timeout}
232 +
233 + client := http.Client{
234 + Timeout: cfg.Timeout,
235 + Transport: &http.Transport{
236 + TLSClientConfig: tlsConfig,
237 + DialContext: d.DialContext,
238 + TLSHandshakeTimeout: cfg.Timeout,
239 + },
240 + }
241 +
242 + return &client, nil
243 +}
244 +
245 +func newTLSConfig(cfg *Config) (*tls.Config, error) {
246 + if cfg.TLS.SrvCertFile == "" && cfg.TLS.SrvKeyFile == "" && cfg.TLS.TrustedCertFile == "" && !cfg.TLS.InsecureSkipVerify {
247 + return nil, nil
248 + }
249 +
250 + var clientCerts []tls.Certificate
251 + if cfg.TLS.SrvCertFile != "" && cfg.TLS.SrvKeyFile != "" {
252 + clientCert, err := tls.LoadX509KeyPair(cfg.TLS.SrvCertFile, cfg.TLS.SrvKeyFile)
253 + if err != nil {
254 + return nil, fmt.Errorf("error loading tls cert and key files: %s", err)
255 + }
256 + clientCerts = append(clientCerts, clientCert)
257 + }
258 +
259 + tlsConfig := &tls.Config{
260 + Certificates: clientCerts,
261 + MinVersion: tls.VersionTLS12,
262 + }
263 +
264 + if cfg.TLS.TrustedCertFile != "" {
265 + caCert, err := os.ReadFile(cfg.TLS.TrustedCertFile)
266 + if err != nil {
267 + return nil, fmt.Errorf("failed to read CA certificate file %q: %w", cfg.TLS.TrustedCertFile, err)
268 + }
269 + caCertPool := x509.NewCertPool()
270 + if !caCertPool.AppendCertsFromPEM(caCert) {
271 + return nil, fmt.Errorf("failed to append CA certificate from %q", cfg.TLS.TrustedCertFile)
272 + }
273 + tlsConfig.RootCAs = caCertPool
274 + tlsConfig.InsecureSkipVerify = false
275 + } else if cfg.TLS.InsecureSkipVerify {
276 + tlsConfig.InsecureSkipVerify = true
277 + }
278 +
279 + return tlsConfig, nil
280 +}
281 +
282 +func closeBody(resp *http.Response) {
283 + if resp != nil && resp.Body != nil {
284 + _, _ = io.Copy(io.Discard, resp.Body)
285 + _ = resp.Body.Close()
286 + }
287 +}
src/go/otel-collector/exporter/journaldexporter/journal_remote_test.go new
+488
@@ -0,0 +1,488 @@
1 +package journaldexporter
2 +
3 +import (
4 + "context"
5 + "errors"
6 + "fmt"
7 + "io"
8 + "net/http"
9 + "net/http/httptest"
10 + "sync"
11 + "testing"
12 + "time"
13 +
14 + "github.com/stretchr/testify/assert"
15 + "github.com/stretchr/testify/require"
16 + "go.uber.org/zap/zaptest"
17 +)
18 +
19 +func TestRemoteJournalClient_sendMessage(t *testing.T) {
20 + tests := map[string]struct {
21 + setupServer func(*testing.T) (*httptest.Server, *chunkWriterHandler)
22 + setupClient func(*testing.T, string) *remoteJournalClient
23 + setupContext func() (context.Context, context.CancelFunc)
24 + messages [][]byte
25 + waitBetween time.Duration
26 + waitAfter time.Duration
27 + validateErr func(*testing.T, error)
28 + validateData func(*testing.T, *chunkWriterHandler, [][]byte)
29 + shutdownFirst bool
30 + }{
31 + "successful message delivery": {
32 + setupServer: prepareTestServer,
33 + setupClient: prepareTestRemoteJournalClient,
34 + setupContext: func() (context.Context, context.CancelFunc) {
35 + return context.Background(), func() {}
36 + },
37 + messages: [][]byte{[]byte("test message data")},
38 + waitAfter: 100 * time.Millisecond,
39 + validateErr: func(t *testing.T, err error) {
40 + assert.NoError(t, err)
41 + },
42 + validateData: func(t *testing.T, handler *chunkWriterHandler, messages [][]byte) {
43 + assert.Equal(t, messages[0], handler.getReceivedData())
44 + },
45 + },
46 + "send empty message": {
47 + setupServer: prepareTestServer,
48 + setupClient: prepareTestRemoteJournalClient,
49 + setupContext: func() (context.Context, context.CancelFunc) {
50 + return context.Background(), func() {}
51 + },
52 + messages: [][]byte{{}},
53 + waitAfter: 100 * time.Millisecond,
54 + validateErr: func(t *testing.T, err error) {
55 + assert.NoError(t, err)
56 + },
57 + validateData: func(t *testing.T, handler *chunkWriterHandler, messages [][]byte) {
58 + assert.Empty(t, handler.getReceivedData())
59 + },
60 + },
61 + "send message with cancelled context": {
62 + setupServer: prepareTestServer,
63 + setupClient: prepareTestRemoteJournalClient,
64 + setupContext: func() (context.Context, context.CancelFunc) {
65 + ctx, cancel := context.WithCancel(context.Background())
66 + cancel() // Cancel immediately
67 + return ctx, cancel
68 + },
69 + messages: [][]byte{[]byte("test message that should not be sent")},
70 + waitAfter: 100 * time.Millisecond,
71 + validateErr: func(t *testing.T, err error) {
72 + assert.Error(t, err)
73 + assert.True(t, errors.Is(err, context.Canceled))
74 + },
75 + },
76 + "attempt to send to non-existent server": {
77 + setupServer: nil, // No server
78 + setupClient: func(t *testing.T, _ string) *remoteJournalClient {
79 + return prepareTestRemoteJournalClient(t, "http://non-existent-server:12345")
80 + },
81 + setupContext: func() (context.Context, context.CancelFunc) {
82 + return context.Background(), func() {}
83 + },
84 + messages: [][]byte{[]byte("test message to non-existent server")},
85 + waitBetween: 100 * time.Millisecond,
86 + validateErr: func(t *testing.T, err error) {
87 + assert.Error(t, err)
88 + //assert.Contains(t, err.Error(), "connection unavailable")
89 + },
90 + },
91 + "send after client is shut down": {
92 + setupServer: prepareTestServer,
93 + setupClient: prepareTestRemoteJournalClient,
94 + setupContext: func() (context.Context, context.CancelFunc) {
95 + return context.Background(), func() {}
96 + },
97 + messages: [][]byte{[]byte("test message after shutdown")},
98 + shutdownFirst: true,
99 + validateErr: func(t *testing.T, err error) {
100 + assert.Error(t, err)
101 + assert.Contains(t, err.Error(), "client is shut down")
102 + },
103 + },
104 + // TODO: fix
105 + //"server responds with error status": {
106 + // setupServer: func(t *testing.T) (*httptest.Server, *chunkWriterHandler) {
107 + // handler := &chunkWriterHandler{
108 + // statusCode: http.StatusInternalServerError,
109 + // }
110 + // server := httptest.NewServer(handler)
111 + // t.Cleanup(func() {
112 + // server.Close()
113 + // })
114 + // return server, handler
115 + // },
116 + // setupClient: prepareTestRemoteJournalClient,
117 + // setupContext: func() (context.Context, context.CancelFunc) {
118 + // return context.Background(), func() {}
119 + // },
120 + // messages: [][]byte{
121 + // []byte("test message expecting error"),
122 + // []byte("second message after error"),
123 + // },
124 + // waitBetween: 300 * time.Millisecond, // Longer wait to ensure background error is processed
125 + // waitAfter: 100 * time.Millisecond,
126 + // validateErr: func(t *testing.T, err error) {
127 + // // The first message might succeed because the error only occurs after the server responds
128 + // // But the second message should fail once the background goroutine sets reqErr
129 + // require.Error(t, err)
130 + // // If we get an error about pipe being closed or connection unavailable, that's valid
131 + // valid := strings.Contains(err.Error(), "connection unavailable") ||
132 + // strings.Contains(err.Error(), "failed to write message") ||
133 + // strings.Contains(err.Error(), "pipe is closed")
134 + // assert.True(t, valid, "Expected error to indicate connection problem, got: %v", err)
135 + // },
136 + //},
137 + // TODO: fix
138 + //"server closes connection mid-stream": {
139 + // setupServer: func(t *testing.T) (*httptest.Server, *chunkWriterHandler) {
140 + // handler := &chunkWriterHandler{
141 + // statusCode: http.StatusOK,
142 + // closeEarly: true,
143 + // }
144 + // server := httptest.NewServer(handler)
145 + // t.Cleanup(func() {
146 + // server.Close()
147 + // })
148 + // return server, handler
149 + // },
150 + // setupClient: prepareTestRemoteJournalClient,
151 + // setupContext: func() (context.Context, context.CancelFunc) {
152 + // return context.Background(), func() {}
153 + // },
154 + // messages: [][]byte{[]byte("first message that will cause server to close"), []byte("second message after closure")},
155 + // waitBetween: 200 * time.Millisecond,
156 + // validateErr: func(t *testing.T, err error) {
157 + // assert.Error(t, err)
158 + // assert.Contains(t, err.Error(), "connection unavailable")
159 + // },
160 + //},
161 + "sending multiple messages": {
162 + setupServer: prepareTestServer,
163 + setupClient: prepareTestRemoteJournalClient,
164 + setupContext: func() (context.Context, context.CancelFunc) {
165 + return context.Background(), func() {}
166 + },
167 + messages: [][]byte{
168 + []byte("first message"),
169 + []byte("second message"),
170 + []byte("third message"),
171 + },
172 + waitAfter: 100 * time.Millisecond,
173 + validateErr: func(t *testing.T, err error) {
174 + assert.NoError(t, err)
175 + },
176 + validateData: func(t *testing.T, handler *chunkWriterHandler, messages [][]byte) {
177 + expected := append(append(messages[0], messages[1]...), messages[2]...)
178 + assert.Equal(t, expected, handler.getReceivedData())
179 + },
180 + },
181 + }
182 +
183 + for name, tc := range tests {
184 + t.Run(name, func(t *testing.T) {
185 + var server *httptest.Server
186 + var handler *chunkWriterHandler
187 + var client *remoteJournalClient
188 +
189 + if tc.setupServer != nil {
190 + server, handler = tc.setupServer(t)
191 + client = tc.setupClient(t, server.URL)
192 + } else {
193 + client = tc.setupClient(t, "")
194 + }
195 + defer func() { _ = client.shutdown(context.Background()) }()
196 +
197 + ctx, cancel := tc.setupContext()
198 + defer cancel()
199 +
200 + if tc.shutdownFirst {
201 + err := client.shutdown(context.Background())
202 + assert.NoError(t, err)
203 + }
204 +
205 + var lastErr error
206 + for i, msg := range tc.messages {
207 + lastErr = client.sendMessage(ctx, msg)
208 +
209 + // For multi-message tests, we only validate error on the last message
210 + if i < len(tc.messages)-1 && tc.waitBetween > 0 {
211 + time.Sleep(tc.waitBetween)
212 + }
213 + }
214 +
215 + if tc.waitAfter > 0 {
216 + time.Sleep(tc.waitAfter)
217 + }
218 +
219 + if tc.validateErr != nil {
220 + tc.validateErr(t, lastErr)
221 + }
222 +
223 + if tc.validateData != nil && handler != nil {
224 + tc.validateData(t, handler, tc.messages)
225 + }
226 + })
227 + }
228 +}
229 +
230 +func TestRemoteJournalClient_shutdown(t *testing.T) {
231 + tests := map[string]struct {
232 + setupServer func(*testing.T) (*httptest.Server, *chunkWriterHandler)
233 + sendMessages func(*testing.T, context.Context, *remoteJournalClient)
234 + setupContext func() (context.Context, context.CancelFunc)
235 + validateErr func(*testing.T, error)
236 + validateShutdown func(*testing.T, *remoteJournalClient, *chunkWriterHandler)
237 + doubleShutdown bool
238 + }{
239 + "normal shutdown": {
240 + setupServer: prepareTestServer,
241 + sendMessages: func(t *testing.T, ctx context.Context, client *remoteJournalClient) {
242 + msg := []byte("test message before shutdown")
243 + err := client.sendMessage(ctx, msg)
244 + assert.NoError(t, err)
245 + },
246 + setupContext: func() (context.Context, context.CancelFunc) {
247 + return context.Background(), func() {}
248 + },
249 + validateErr: func(t *testing.T, err error) {
250 + assert.NoError(t, err)
251 + },
252 + validateShutdown: func(t *testing.T, client *remoteJournalClient, handler *chunkWriterHandler) {
253 + select {
254 + case <-client.done:
255 + // Expected - channel should be closed
256 + default:
257 + t.Error("Expected done channel to be closed after shutdown")
258 + }
259 +
260 + err := client.sendMessage(context.Background(), []byte("test message after shutdown"))
261 + assert.Error(t, err)
262 + assert.Contains(t, err.Error(), "client is shut down")
263 + },
264 + },
265 + // TODO: fix
266 + //"shutdown with timeout": {
267 + // setupServer: func(t *testing.T) (*httptest.Server, *chunkWriterHandler) {
268 + // handler := &chunkWriterHandler{
269 + // statusCode: http.StatusOK,
270 + // delay: 500 * time.Millisecond, // Delay each read to simulate slow processing
271 + // }
272 + // server := httptest.NewServer(handler)
273 + // t.Cleanup(func() {
274 + // server.Close()
275 + // })
276 + // return server, handler
277 + // },
278 + // sendMessages: func(t *testing.T, ctx context.Context, client *remoteJournalClient) {
279 + // // Send a large message to ensure the background task will be busy
280 + // msg := []byte(strings.Repeat("test message with shutdown timeout ", 1000))
281 + // err := client.sendMessage(ctx, msg)
282 + // assert.NoError(t, err)
283 + // },
284 + // setupContext: func() (context.Context, context.CancelFunc) {
285 + // return context.WithTimeout(context.Background(), 100*time.Millisecond)
286 + // },
287 + // validateErr: func(t *testing.T, err error) {
288 + // assert.Error(t, err)
289 + // assert.True(t, errors.Is(err, context.DeadlineExceeded))
290 + // },
291 + // validateShutdown: func(t *testing.T, client *remoteJournalClient, handler *chunkWriterHandler) {
292 + // select {
293 + // case <-client.done:
294 + // // Expected - channel should be closed
295 + // default:
296 + // t.Error("Expected done channel to be closed after shutdown timeout")
297 + // }
298 + // },
299 + //},
300 + "double shutdown": {
301 + setupServer: prepareTestServer,
302 + sendMessages: func(t *testing.T, ctx context.Context, client *remoteJournalClient) {
303 + // No messages sent
304 + },
305 + setupContext: func() (context.Context, context.CancelFunc) {
306 + return context.Background(), func() {}
307 + },
308 + validateErr: func(t *testing.T, err error) {
309 + assert.NoError(t, err)
310 + },
311 + doubleShutdown: true,
312 + },
313 + "shutdown with active connections": {
314 + setupServer: prepareTestServer,
315 + sendMessages: func(t *testing.T, ctx context.Context, client *remoteJournalClient) {
316 + // Send multiple messages to ensure there's an active connection
317 + for i := 0; i < 5; i++ {
318 + msg := []byte(fmt.Sprintf("test message %d before shutdown", i))
319 + err := client.sendMessage(ctx, msg)
320 + assert.NoError(t, err)
321 + }
322 + // Allow some time for the background goroutine to process
323 + time.Sleep(50 * time.Millisecond)
324 + },
325 + setupContext: func() (context.Context, context.CancelFunc) {
326 + return context.Background(), func() {}
327 + },
328 + validateErr: func(t *testing.T, err error) {
329 + assert.NoError(t, err)
330 + },
331 + validateShutdown: func(t *testing.T, client *remoteJournalClient, handler *chunkWriterHandler) {
332 + // Verify at least some data was received before shutdown
333 + receivedData := handler.getReceivedData()
334 + assert.NotEmpty(t, receivedData)
335 + assert.Contains(t, string(receivedData), "test message")
336 + },
337 + },
338 + "shutdown without active connections": {
339 + setupServer: prepareTestServer,
340 + sendMessages: func(t *testing.T, ctx context.Context, client *remoteJournalClient) {
341 + // No messages sent
342 + },
343 + setupContext: func() (context.Context, context.CancelFunc) {
344 + return context.Background(), func() {}
345 + },
346 + validateErr: func(t *testing.T, err error) {
347 + assert.NoError(t, err)
348 + },
349 + validateShutdown: func(t *testing.T, client *remoteJournalClient, handler *chunkWriterHandler) {
350 + select {
351 + case <-client.done:
352 + // Expected - channel should be closed
353 + default:
354 + t.Error("Expected done channel to be closed after shutdown")
355 + }
356 + },
357 + },
358 + }
359 +
360 + for name, tc := range tests {
361 + t.Run(name, func(t *testing.T) {
362 + server, handler := tc.setupServer(t)
363 +
364 + client := prepareTestRemoteJournalClient(t, server.URL)
365 +
366 + // Setup context for messages and initial operations
367 + msgCtx, msgCancel := context.WithTimeout(context.Background(), 5*time.Second)
368 + defer msgCancel()
369 +
370 + tc.sendMessages(t, msgCtx, client)
371 +
372 + shutdownCtx, shutdownCancel := tc.setupContext()
373 + defer shutdownCancel()
374 +
375 + err := client.shutdown(shutdownCtx)
376 +
377 + if tc.doubleShutdown {
378 + secondErr := client.shutdown(context.Background())
379 + assert.NoError(t, secondErr, "Second shutdown should not error")
380 + }
381 +
382 + if tc.validateErr != nil {
383 + tc.validateErr(t, err)
384 + }
385 +
386 + if tc.validateShutdown != nil {
387 + tc.validateShutdown(t, client, handler)
388 + }
389 + })
390 + }
391 +}
392 +
393 +type chunkWriterHandler struct {
394 + receivedData []byte
395 + mu sync.Mutex
396 + delay time.Duration
397 + statusCode int
398 + errorAfterRead bool
399 + closeEarly bool
400 +}
401 +
402 +func (h *chunkWriterHandler) ServeHTTP(w http.ResponseWriter, r *http.Request) {
403 + h.mu.Lock()
404 + h.receivedData = make([]byte, 0)
405 + h.mu.Unlock()
406 +
407 + contentType := r.Header.Get("Content-Type")
408 + if contentType != "application/vnd.fdo.journal" {
409 + w.WriteHeader(http.StatusBadRequest)
410 + _, _ = w.Write([]byte("Invalid content type"))
411 + return
412 + }
413 +
414 + w.Header().Set("Transfer-Encoding", "chunked")
415 + w.WriteHeader(h.statusCode)
416 +
417 + buf := make([]byte, 1024)
418 + readCount := 0
419 +
420 + readDone := make(chan struct{})
421 + go func() {
422 + defer close(readDone)
423 +
424 + for {
425 + if h.delay > 0 {
426 + time.Sleep(h.delay)
427 + }
428 + n, err := r.Body.Read(buf)
429 + if n > 0 {
430 + h.mu.Lock()
431 + h.receivedData = append(h.receivedData, buf[:n]...)
432 + h.mu.Unlock()
433 + readCount++
434 + }
435 + if err == io.EOF {
436 + break
437 + }
438 + if err != nil {
439 + return
440 + }
441 + if h.errorAfterRead && readCount >= 1 {
442 + if flusher, ok := w.(http.Flusher); ok {
443 + flusher.Flush()
444 + }
445 + return
446 + }
447 + if h.closeEarly && readCount >= 1 {
448 + return
449 + }
450 + if h.statusCode < 200 || h.statusCode >= 300 {
451 + return
452 + }
453 + }
454 + }()
455 +
456 + <-readDone
457 +}
458 +
459 +func (h *chunkWriterHandler) getReceivedData() []byte {
460 + h.mu.Lock()
461 + defer h.mu.Unlock()
462 + return h.receivedData
463 +}
464 +
465 +func prepareTestServer(t *testing.T) (*httptest.Server, *chunkWriterHandler) {
466 + handler := &chunkWriterHandler{
467 + statusCode: http.StatusOK,
468 + }
469 + server := httptest.NewServer(handler)
470 + t.Cleanup(func() {
471 + server.Close()
472 + })
473 + return server, handler
474 +}
475 +
476 +func prepareTestRemoteJournalClient(t *testing.T, serverURL string) *remoteJournalClient {
477 + cfg := &Config{
478 + URL: serverURL,
479 + Timeout: 1 * time.Second,
480 + }
481 +
482 + client, err := newRemoteJournalClient(cfg)
483 + require.NoError(t, err)
484 + require.NotNil(t, client)
485 +
486 + client.log = zaptest.NewLogger(t)
487 + return client
488 +}
src/go/otel-collector/exporter/journaldexporter/journal_socket.go new
+190
@@ -0,0 +1,190 @@
1 +package journaldexporter
2 +
3 +import (
4 + "context"
5 + "errors"
6 + "fmt"
7 + "net"
8 + "os"
9 + "sync"
10 + "syscall"
11 + "time"
12 +)
13 +
14 +const defaultJournalSocket = "/run/systemd/journal/socket"
15 +
16 +func newSocketJournalClient() (*socketJournalClient, error) {
17 + if _, err := os.Stat(defaultJournalSocket); os.IsNotExist(err) {
18 + return nil, fmt.Errorf("journal socket does not exist: %w", err)
19 + }
20 +
21 + conn, err := createJournalConn()
22 + if err != nil {
23 + return nil, fmt.Errorf("failed to create journal connection: %w", err)
24 + }
25 +
26 + tempFile, err := createUnlinkedTempFile()
27 + if err != nil {
28 + _ = conn.Close()
29 + return nil, fmt.Errorf("failed to create temporary file: %w", err)
30 + }
31 +
32 + return &socketJournalClient{
33 + socketPath: defaultJournalSocket,
34 + conn: conn,
35 + done: make(chan struct{}),
36 + tempFile: tempFile,
37 + }, nil
38 +}
39 +
40 +type socketJournalClient struct {
41 + socketPath string
42 + conn *net.UnixConn
43 + mu sync.Mutex
44 + done chan struct{}
45 +
46 + // Pre-created temporary file for large messages
47 + tempFile *os.File
48 +}
49 +
50 +func (jc *socketJournalClient) sendMessage(ctx context.Context, msg []byte) error {
51 + if len(msg) == 0 {
52 + return nil
53 + }
54 + if jc.conn == nil {
55 + return errors.New("journal: connection is closed")
56 + }
57 +
58 + jc.mu.Lock()
59 + defer jc.mu.Unlock()
60 +
61 + select {
62 + case <-ctx.Done():
63 + return ctx.Err()
64 + case <-jc.done:
65 + return errors.New("journal: client is shut down")
66 + default:
67 + }
68 +
69 + socketAddr := &net.UnixAddr{
70 + Name: jc.socketPath,
71 + Net: "unixgram",
72 + }
73 +
74 + if err := jc.setConnWriteDeadline(ctx); err != nil {
75 + return fmt.Errorf("journal: failed to set write deadline: %w", err)
76 + }
77 +
78 + if _, _, err := jc.conn.WriteMsgUnix(msg, nil, socketAddr); err != nil {
79 + if !isSocketSpaceError(err) {
80 + return fmt.Errorf("journal: failed to write to socket: %w", err)
81 + }
82 + return jc.sendViaFd(ctx, msg, socketAddr)
83 + }
84 +
85 + return nil
86 +}
87 +
88 +func (jc *socketJournalClient) sendViaFd(ctx context.Context, msg []byte, socketAddr *net.UnixAddr) error {
89 + if _, err := jc.tempFile.Seek(0, 0); err != nil {
90 + return fmt.Errorf("journal: failed to seek in temporary file: %w", err)
91 + }
92 +
93 + if err := jc.tempFile.Truncate(0); err != nil {
94 + return fmt.Errorf("journal: failed to truncate temporary file: %w", err)
95 + }
96 +
97 + if _, err := jc.tempFile.Write(msg); err != nil {
98 + return fmt.Errorf("journal: failed to write to temporary file: %w", err)
99 + }
100 +
101 + if _, err := jc.tempFile.Seek(0, 0); err != nil {
102 + return fmt.Errorf("journal: failed to reset file position: %w", err)
103 + }
104 +
105 + rights := syscall.UnixRights(int(jc.tempFile.Fd()))
106 +
107 + if err := jc.setConnWriteDeadline(ctx); err != nil {
108 + return fmt.Errorf("journal: failed to set write deadline: %w", err)
109 + }
110 +
111 + if _, _, err := jc.conn.WriteMsgUnix([]byte{}, rights, socketAddr); err != nil {
112 + return fmt.Errorf("journal: failed to send file descriptor: %w", err)
113 + }
114 +
115 + return nil
116 +}
117 +
118 +func (jc *socketJournalClient) shutdown(ctx context.Context) error {
119 + jc.mu.Lock()
120 + defer jc.mu.Unlock()
121 +
122 + select {
123 + case <-jc.done:
124 + return nil
125 + default:
126 + close(jc.done)
127 + }
128 +
129 + if jc.tempFile != nil {
130 + _ = jc.tempFile.Close()
131 + jc.tempFile = nil
132 + }
133 +
134 + if jc.conn != nil {
135 + _ = jc.conn.Close()
136 + jc.conn = nil
137 + }
138 +
139 + return nil
140 +}
141 +
142 +func (jc *socketJournalClient) setConnWriteDeadline(ctx context.Context) error {
143 + var timeout = 5 * time.Second
144 + if deadline, ok := ctx.Deadline(); ok {
145 + timeout = time.Until(deadline)
146 + }
147 + return jc.conn.SetWriteDeadline(time.Now().Add(timeout))
148 +}
149 +
150 +func createJournalConn() (*net.UnixConn, error) {
151 + autobind, err := net.ResolveUnixAddr("unixgram", "")
152 + if err != nil {
153 + return nil, fmt.Errorf("failed to resolve unix address: %w", err)
154 + }
155 +
156 + conn, err := net.ListenUnixgram("unixgram", autobind)
157 + if err != nil {
158 + return nil, fmt.Errorf("failed to create unix datagram socket: %w", err)
159 + }
160 +
161 + return conn, nil
162 +}
163 +
164 +func createUnlinkedTempFile() (*os.File, error) {
165 + file, err := os.CreateTemp("/dev/shm/", "journal.XXXXX")
166 + if err != nil {
167 + return nil, err
168 + }
169 +
170 + // Unlink the file so it's automatically cleaned up when closed
171 + err = syscall.Unlink(file.Name())
172 + if err != nil {
173 + _ = file.Close()
174 + return nil, err
175 + }
176 +
177 + return file, nil
178 +}
179 +
180 +func isSocketSpaceError(err error) bool {
181 + // checks whether the error is signaling an "overlarge message" condition
182 + var opErr *net.OpError
183 + var sysErr *os.SyscallError
184 +
185 + if !errors.As(err, &opErr) || !errors.As(opErr.Err, &sysErr) {
186 + return false
187 + }
188 +
189 + return errors.Is(sysErr.Err, syscall.EMSGSIZE) || errors.Is(sysErr.Err, syscall.ENOBUFS)
190 +}