1 package reaper
2
3 import (
4 "context"
5 "os"
6 "path/filepath"
7 "regexp"
8 "sync"
9 "testing"
10 "time"
11
12 "github.com/cybertec-postgresql/pgwatch/v6/internal/metrics"
13 "github.com/cybertec-postgresql/pgwatch/v6/internal/sources"
14 "github.com/cybertec-postgresql/pgwatch/v6/internal/testutil"
15 "github.com/pashagolub/pgxmock/v5"
16 "github.com/stretchr/testify/assert"
17 "github.com/stretchr/testify/require"
18 )
19
20 var expectedSettingsQuery = `select current_setting`
21
22 func TestNewLogParser(t *testing.T) {
23 tempDir := t.TempDir()
24
25 mock, err := pgxmock.NewPool()
26 require.NoError(t, err)
27 defer mock.Close()
28
29 sourceConn := &sources.DbConn{
30 Source: sources.Source{
31 Name: "test-source",
32 Metrics: metrics.MetricIntervals{specialMetricServerLogEventCounts: 60.0},
33 },
34 Conn: mock,
35 }
36 storeCh := make(chan metrics.MeasurementEnvelope, 10)
37
38 t.Run("success", func(t *testing.T) {
39 mock.ExpectQuery(expectedSettingsQuery).
40 WillReturnRows(pgxmock.NewRows([]string{"is_enabled", "csvlog_dest", "log_trunc", "log_dir", "lc_messages"}).
41 AddRow(true, true, false, tempDir, "en"))
42
43 lp, err := NewLogParser(testutil.TestContext, sourceConn, storeCh)
44
45 assert.NoError(t, err)
46 assert.NotNil(t, lp)
47 assert.Equal(t, true, lp.CollectorEnabled)
48 assert.True(t, lp.CSVDestination)
49 assert.Equal(t, tempDir, lp.Directory)
50 assert.Equal(t, "en", lp.ServerMessagesLang)
51 assert.Equal(t, false, lp.TruncateOnRotation)
52 assert.Equal(t, 60*time.Second, lp.Interval)
53 assert.NotNil(t, lp.LogsMatchRegex)
54 assert.NoError(t, mock.ExpectationsWereMet())
55 })
56
57 t.Run("tryDetermineLogSettings error", func(t *testing.T) {
58 mock.ExpectQuery(expectedSettingsQuery).WillReturnError(assert.AnError)
59 lp, err := NewLogParser(testutil.TestContext, sourceConn, storeCh)
60 assert.Error(t, err)
61 assert.Nil(t, lp)
62 assert.NoError(t, mock.ExpectationsWereMet())
63 })
64
65 t.Run("unknown language defaults to en", func(t *testing.T) {
66 mock.ExpectQuery(expectedSettingsQuery).
67 WillReturnRows(pgxmock.NewRows([]string{"is_enabled", "csvlog_dest", "log_trunc", "log_dir", "lc_messages"}).
68 AddRow(true, true, false, tempDir, "zz"))
69
70 lp, err := NewLogParser(testutil.TestContext, sourceConn, storeCh)
71 assert.NoError(t, err)
72 assert.NotNil(t, lp)
73 assert.Equal(t, "en", lp.ServerMessagesLang)
74 assert.NoError(t, mock.ExpectationsWereMet())
75 })
76
77 t.Run("relative log directory", func(t *testing.T) {
78 mock.ExpectQuery(expectedSettingsQuery).
79 WillReturnRows(pgxmock.NewRows([]string{"is_enabled", "csvlog_dest", "log_trunc", "log_dir", "lc_messages"}).
80 AddRow(true, true, true, "/data/pg_log", "de"))
81
82 lp, err := NewLogParser(testutil.TestContext, sourceConn, storeCh)
83 assert.NoError(t, err)
84 assert.NotNil(t, lp)
85 assert.Equal(t, "/data/pg_log", lp.Directory)
86 assert.Equal(t, "de", lp.ServerMessagesLang)
87 assert.Equal(t, true, lp.TruncateOnRotation)
88 assert.NoError(t, mock.ExpectationsWereMet())
89 })
90
91 t.Run("logging_collector disabled", func(t *testing.T) {
92 mock.ExpectQuery(expectedSettingsQuery).
93 WillReturnRows(pgxmock.NewRows([]string{"is_enabled", "csvlog_dest", "log_trunc", "log_dir", "lc_messages"}).
94 AddRow(false, true, true, "/data/pg_log", "de"))
95
96 lp, err := NewLogParser(testutil.TestContext, sourceConn, storeCh)
97 assert.Equal(t, err.Error(), "logging_collector is not enabled on the db server")
98 assert.Nil(t, lp)
99 })
100
101 t.Run("csvlog not in log_destination", func(t *testing.T) {
102 mock.ExpectQuery(expectedSettingsQuery).
103 WillReturnRows(pgxmock.NewRows([]string{"is_enabled", "csvlog_dest", "log_trunc", "log_dir", "lc_messages"}).
104 AddRow(true, false, true, "/data/pg_log", "de"))
105
106 lp, err := NewLogParser(testutil.TestContext, sourceConn, storeCh)
107 assert.Equal(t, err.Error(), "log_destination must contain 'csvlog' for log parsing to work")
108 assert.Nil(t, lp)
109 })
110 }
111
112 func TestTryDetermineLogSettings(t *testing.T) {
113 t.Run("absolute log directory - known lang", func(t *testing.T) {
114 mock, err := pgxmock.NewPool()
115 require.NoError(t, err)
116 defer mock.Close()
117
118 mock.ExpectQuery(expectedSettingsQuery).
119 WillReturnRows(pgxmock.NewRows([]string{"is_enabled", "csvlog_dest", "log_trunc", "log_dir", "lc_messages"}).
120 AddRow(true, true, false, "/var/log/postgresql", "de"))
121
122 logCfg, err := tryDetermineLogSettings(testutil.TestContext, mock)
123 assert.NoError(t, err)
124 assert.Equal(t, "/var/log/postgresql", logCfg.Directory)
125 assert.Equal(t, "de", logCfg.ServerMessagesLang)
126 assert.Equal(t, false, logCfg.TruncateOnRotation)
127 assert.NoError(t, mock.ExpectationsWereMet())
128 })
129
130 t.Run("relative log directory - unknown lang", func(t *testing.T) {
131 mock, err := pgxmock.NewPool()
132 require.NoError(t, err)
133 defer mock.Close()
134
135 mock.ExpectQuery(expectedSettingsQuery).
136 WillReturnRows(pgxmock.NewRows([]string{"is_enabled", "csvlog_dest", "log_trunc", "log_dir", "lc_messages"}).
137 AddRow(true, true, false, "/data/log", "xx"))
138
139 logCfg, err := tryDetermineLogSettings(testutil.TestContext, mock)
140 assert.NoError(t, err)
141 assert.Equal(t, "/data/log", logCfg.Directory)
142 assert.Equal(t, "en", logCfg.ServerMessagesLang)
143 assert.Equal(t, false, logCfg.TruncateOnRotation)
144 assert.NoError(t, mock.ExpectationsWereMet())
145 })
146
147 t.Run("query error", func(t *testing.T) {
148 mock, err := pgxmock.NewPool()
149 require.NoError(t, err)
150 defer mock.Close()
151
152 mock.ExpectQuery(expectedSettingsQuery).
153 WillReturnError(assert.AnError)
154
155 logCfg, err := tryDetermineLogSettings(testutil.TestContext, mock)
156 assert.Error(t, err)
157 assert.Nil(t, logCfg)
158 assert.NoError(t, mock.ExpectationsWereMet())
159 })
160 }
161
162 func TestCheckHasPrivileges(t *testing.T) {
163 tempDir := t.TempDir()
164
165 names := [2]string{"pg_ls_logdir() fails", "pg_read_file() permission denied"}
166 for _, name := range names {
167 t.Run("checkHasRemotePrivileges fails - "+name, func(t *testing.T) {
168 mock, err := pgxmock.NewPool()
169 require.NoError(t, err)
170 defer mock.Close()
171
172 mock.ExpectQuery(expectedSettingsQuery).
173 WillReturnRows(pgxmock.NewRows([]string{"is_enabled", "csvlog_dest", "log_trunc", "log_dir", "lc_messages"}).
174 AddRow(true, true, false, tempDir, "en"))
175
176
177 mock.ExpectQuery(`SELECT COALESCE`).WillReturnRows(
178 pgxmock.NewRows([]string{"is_unix_socket"}).AddRow(false))
179
180 if name == "pg_ls_logdir() fails" {
181
182 mock.ExpectQuery(`select name from pg_ls_logdir\(\) limit 1`).
183 WillReturnError(assert.AnError)
184 } else {
185
186 mock.ExpectQuery(`select name from pg_ls_logdir\(\) limit 1`).
187 WillReturnRows(pgxmock.NewRows([]string{"name"}).AddRow("log.csv"))
188
189
190 mock.ExpectQuery(`select pg_read_file\(\$1, 0, 0\)`).
191 WithArgs(filepath.Join(tempDir, "log.csv")).
192 WillReturnError(assert.AnError)
193 }
194
195 sourceConn := &sources.DbConn{
196 Source: sources.Source{
197 Name: "test-source",
198 },
199 Conn: mock,
200 }
201
202 storeCh := make(chan metrics.MeasurementEnvelope, 10)
203
204 lp, err := NewLogParser(testutil.TestContext, sourceConn, storeCh)
205 require.NoError(t, err)
206
207 err = lp.ParseLogs()
208 assert.Error(t, err)
209
210
211 assert.NoError(t, mock.ExpectationsWereMet())
212
213
214 select {
215 case measurement := <-storeCh:
216 t.Errorf("Expected no data, but got: %+v", measurement)
217 case <-time.After(time.Second):
218
219 }
220 })
221 }
222 }
223
224 func TestEventCountsToMetricStoreMessages(t *testing.T) {
225 mdb := &sources.DbConn{
226 Source: sources.Source{
227 Name: "test-db",
228 Kind: sources.SourcePostgres,
229 CustomTags: map[string]string{"env": "test"},
230 },
231 }
232 lp := &LogParser{
233 SourceConn: mdb,
234 eventCounts: map[string]int64{
235 "ERROR": 5,
236 "WARNING": 10,
237 },
238 eventCountsTotal: map[string]int64{
239 "ERROR": 15,
240 "WARNING": 25,
241 "INFO": 50,
242 },
243 }
244 result := lp.GetMeasurementEnvelope()
245
246 assert.Equal(t, "test-db", result.DBName)
247 assert.Equal(t, specialMetricServerLogEventCounts, result.MetricName)
248 assert.Equal(t, map[string]string{"env": "test"}, result.CustomTags)
249
250
251 assert.Len(t, result.Data, 1)
252 measurement := result.Data[0]
253
254
255 assert.Equal(t, int64(5), measurement["error"])
256 assert.Equal(t, int64(10), measurement["warning"])
257 assert.Equal(t, int64(0), measurement["info"])
258 assert.Equal(t, int64(0), measurement["debug"])
259
260
261 assert.Equal(t, int64(15), measurement["error_total"])
262 assert.Equal(t, int64(25), measurement["warning_total"])
263 assert.Equal(t, int64(50), measurement["info_total"])
264 assert.Equal(t, int64(0), measurement["debug_total"])
265 }
266
267 func TestSeverityToEnglish(t *testing.T) {
268 tests := []struct {
269 serverLang string
270 errorSeverity string
271 expected string
272 }{
273 {"en", "ERROR", "ERROR"},
274 {"de", "FEHLER", "ERROR"},
275 {"fr", "ERREUR", "ERROR"},
276 {"de", "WARNUNG", "WARNING"},
277 {"ru", "ОШИБКА", "ERROR"},
278 {"zh", "错误", "ERROR"},
279 {"unknown", "ERROR", "ERROR"},
280 {"de", "UNKNOWN_SEVERITY", "UNKNOWN_SEVERITY"},
281 }
282
283 for _, tt := range tests {
284 t.Run(tt.serverLang+"_"+tt.errorSeverity, func(t *testing.T) {
285 result := severityToEnglish(tt.serverLang, tt.errorSeverity)
286 assert.Equal(t, tt.expected, result)
287 })
288 }
289 }
290
291 func TestZeroEventCounts(t *testing.T) {
292 eventCounts := map[string]int64{
293 "ERROR": 5,
294 "WARNING": 10,
295 "INFO": 15,
296 }
297
298 zeroEventCounts(eventCounts)
299
300
301 for _, severity := range pgSeverities {
302 assert.Equal(t, int64(0), eventCounts[severity])
303 }
304 }
305
306 func TestRegexMatchesToMap(t *testing.T) {
307 t.Run("successful match", func(t *testing.T) {
308 lp := &LogParser{
309 LogsMatchRegex: regexp.MustCompile(`(?P<severity>\w+): (?P<message>.+)`),
310 }
311 matches := []string{"ERROR: Something went wrong", "ERROR", "Something went wrong"}
312
313 result := lp.regexMatchesToMap(matches)
314 expected := map[string]string{
315 "severity": "ERROR",
316 "message": "Something went wrong",
317 }
318
319 assert.Equal(t, expected, result)
320 })
321
322 t.Run("no matches", func(t *testing.T) {
323 lp := &LogParser{
324 LogsMatchRegex: regexp.MustCompile(`(?P<severity>\w+): (?P<message>.+)`),
325 }
326 matches := []string{}
327
328 result := lp.regexMatchesToMap(matches)
329 assert.Empty(t, result)
330 })
331
332 t.Run("nil regex", func(t *testing.T) {
333 lp := &LogParser{}
334 matches := []string{"test"}
335
336 result := lp.regexMatchesToMap(matches)
337 assert.Empty(t, result)
338 })
339 }
340
341 func TestCSVLogRegex(t *testing.T) {
342
343 lp := &LogParser{
344 LogsMatchRegex: regexp.MustCompile(csvLogDefaultRegEx),
345 }
346
347 testLines := []struct {
348 line string
349 expected map[string]string
350 }{
351 {
352 line: `2023-12-01 10:30:45.123 UTC,"postgres","testdb",12345,"127.0.0.1:54321",session123,1,"SELECT",2023-12-01 10:30:00 UTC,1/234,567,ERROR,`,
353 expected: map[string]string{
354 "log_time": "2023-12-01 10:30:45.123 UTC",
355 "user_name": "postgres",
356 "database_name": "testdb",
357 "process_id": "12345",
358 "connection_from": "127.0.0.1:54321",
359 "session_id": "session123",
360 "session_line_num": "1",
361 "command_tag": "SELECT",
362 "error_severity": "ERROR",
363 },
364 },
365 {
366 line: `2023-12-01 10:30:45.123 UTC,postgres,testdb,12345,127.0.0.1:54321,session123,1,SELECT,2023-12-01 10:30:00 UTC,1/234,567,WARNING,`,
367 expected: map[string]string{
368 "log_time": "2023-12-01 10:30:45.123 UTC",
369 "user_name": "postgres",
370 "database_name": "testdb",
371 "process_id": "12345",
372 "connection_from": "127.0.0.1:54321",
373 "session_id": "session123",
374 "session_line_num": "1",
375 "command_tag": "SELECT",
376 "error_severity": "WARNING",
377 },
378 },
379 }
380
381 for i, tt := range testLines {
382 t.Run(string(rune('A'+i)), func(t *testing.T) {
383 matches := lp.LogsMatchRegex.FindStringSubmatch(tt.line)
384 assert.NotEmpty(t, matches, "regex should match the log line")
385
386 result := lp.regexMatchesToMap(matches)
387 for key, expected := range tt.expected {
388 assert.Equal(t, expected, result[key], "mismatch for key %s", key)
389 }
390 })
391 }
392 }
393
394 func TestLogParseLocal(t *testing.T) {
395 tempDir := t.TempDir()
396 logFile := filepath.Join(tempDir, "test.csv")
397
398
399 logContent := `2023-12-01 10:30:45.123 UTC,"postgres","testdb",12345,"127.0.0.1:54321",session123,1,"SELECT",2023-12-01 10:30:00 UTC,1/234,567,ERROR,"duplicate key value violates unique constraint"
400 2023-12-01 10:30:46.124 UTC,"postgres","testdb",12345,"127.0.0.1:54321",session123,2,"SELECT",2023-12-01 10:30:00 UTC,1/234,567,WARNING,"this is a warning message"
401 2023-12-01 10:30:47.125 UTC,"postgres","otherdb",12346,"127.0.0.1:54322",session124,1,"INSERT",2023-12-01 10:30:00 UTC,1/235,568,ERROR,"another error message"
402 `
403
404 err := os.WriteFile(logFile, []byte(logContent), 0644)
405 require.NoError(t, err)
406
407
408 mock, err := pgxmock.NewPool()
409 require.NoError(t, err)
410 defer mock.Close()
411
412 mock.ExpectQuery(expectedSettingsQuery).
413 WillReturnRows(pgxmock.NewRows([]string{"is_enabled", "csvlog_dest", "log_trunc", "log_dir", "lc_messages"}).
414 AddRow(true, true, false, tempDir, "en"))
415
416 mock.ExpectQuery(`SELECT COALESCE`).WillReturnRows(
417 pgxmock.NewRows([]string{"is_unix_socket"}).AddRow(true))
418
419
420 sourceConn := &sources.DbConn{
421 Source: sources.Source{
422 Name: "test-source",
423 },
424 Conn: mock,
425 }
426
427
428 ctx, cancel := context.WithTimeout(testutil.TestContext, 2*time.Second)
429 defer cancel()
430
431
432 storeCh := make(chan metrics.MeasurementEnvelope, 10)
433
434 lp, err := NewLogParser(ctx, sourceConn, storeCh)
435 require.NoError(t, err)
436 err = lp.ParseLogs()
437 assert.NoError(t, err)
438
439
440 assert.NoError(t, mock.ExpectationsWereMet())
441
442
443 var measurement metrics.MeasurementEnvelope
444 select {
445 case measurement = <-storeCh:
446 assert.NotEmpty(t, measurement.Data, "Measurement data should not be empty")
447 case <-time.After(2 * time.Second):
448 }
449
450 assert.Equal(t, "test-source", measurement.DBName)
451 assert.Equal(t, specialMetricServerLogEventCounts, measurement.MetricName)
452
453
454 data := measurement.Data[0]
455
456 _, hasError := data["error"]
457 _, hasWarning := data["warning"]
458 assert.True(t, hasError && hasWarning, "Should have at least error and warning")
459 }
460
461 func TestGetFileWithLatestTimestamp(t *testing.T) {
462
463 tempDir := t.TempDir()
464
465 t.Run("single file", func(t *testing.T) {
466 file1 := filepath.Join(tempDir, "test1.log")
467 err := os.WriteFile(file1, []byte("test"), 0644)
468 require.NoError(t, err)
469
470 latest, err := getFileWithLatestTimestamp([]string{file1})
471 assert.NoError(t, err)
472 assert.Equal(t, file1, latest)
473 })
474
475 t.Run("multiple files with different timestamps", func(t *testing.T) {
476 file1 := filepath.Join(tempDir, "old.log")
477 file2 := filepath.Join(tempDir, "new.log")
478
479
480 err := os.WriteFile(file1, []byte("old"), 0644)
481 require.NoError(t, err)
482
483
484 time.Sleep(10 * time.Millisecond)
485
486
487 err = os.WriteFile(file2, []byte("new"), 0644)
488 require.NoError(t, err)
489
490 latest, err := getFileWithLatestTimestamp([]string{file1, file2})
491 assert.NoError(t, err)
492 assert.Equal(t, file2, latest)
493 })
494
495 t.Run("empty file list", func(t *testing.T) {
496 latest, err := getFileWithLatestTimestamp([]string{})
497 assert.NoError(t, err)
498 assert.Equal(t, "", latest)
499 })
500
501 t.Run("non-existent file", func(t *testing.T) {
502 nonExistent := filepath.Join(tempDir, "nonexistent.log")
503 latest, err := getFileWithLatestTimestamp([]string{nonExistent})
504 assert.Error(t, err)
505 assert.Equal(t, "", latest)
506 })
507 }
508
509 func TestGetFileWithNextModTimestamp(t *testing.T) {
510 tempDir := t.TempDir()
511
512 t.Run("finds next file", func(t *testing.T) {
513 file1 := filepath.Join(tempDir, "first.log")
514 file2 := filepath.Join(tempDir, "second.log")
515 file3 := filepath.Join(tempDir, "third.log")
516
517
518 err := os.WriteFile(file1, []byte("first"), 0644)
519 require.NoError(t, err)
520
521 time.Sleep(10 * time.Millisecond)
522 err = os.WriteFile(file2, []byte("second"), 0644)
523 require.NoError(t, err)
524
525 time.Sleep(10 * time.Millisecond)
526 err = os.WriteFile(file3, []byte("third"), 0644)
527 require.NoError(t, err)
528
529 globPattern := filepath.Join(tempDir, "*.log")
530 next, err := getFileWithNextModTimestamp(globPattern, file1)
531 assert.NoError(t, err)
532 assert.Equal(t, file2, next)
533 })
534
535 t.Run("no next file", func(t *testing.T) {
536 file1 := filepath.Join(tempDir, "only.log")
537 err := os.WriteFile(file1, []byte("only"), 0644)
538 require.NoError(t, err)
539
540 globPattern := filepath.Join(tempDir, "*.log")
541 next, err := getFileWithNextModTimestamp(globPattern, file1)
542 assert.NoError(t, err)
543 assert.Equal(t, "", next)
544 })
545
546 t.Run("invalid glob pattern", func(t *testing.T) {
547 invalidGlob := "["
548 file1 := filepath.Join(tempDir, "test.log")
549 next, err := getFileWithNextModTimestamp(invalidGlob, file1)
550 assert.Error(t, err)
551 assert.Equal(t, "", next)
552 })
553 }
554
555 func TestLogParseRemote(t *testing.T) {
556 const (
557 testTimeout = 3 * time.Second
558 channelBufferSize = 10
559 logFileName = "postgresql.csv"
560 testDbName = "testdb"
561 )
562
563
564 logContent := `2023-12-01 10:30:45.123 UTC,"postgres","testdb",12345,"127.0.0.1:54321",session123,1,"SELECT",2023-12-01 10:30:00 UTC,1/234,567,ERROR,"duplicate key value violates unique constraint"
565 2023-12-01 10:30:46.124 UTC,"postgres","testdb",12345,"127.0.0.1:54321",session123,2,"SELECT",2023-12-01 10:30:00 UTC,1/234,567,WARNING,"this is a warning message"
566 2023-12-01 10:30:47.125 UTC,"postgres","otherdb",12346,"127.0.0.1:54322",session124,1,"INSERT",2023-12-01 10:30:00 UTC,1/235,568,ERROR,"another error message"
567 `
568
569 t.Run("success - parses CSV logs with correct counts", func(t *testing.T) {
570 tempDir := t.TempDir()
571
572 mock, err := pgxmock.NewPool()
573 require.NoError(t, err)
574 defer mock.Close()
575
576 mock.ExpectQuery(expectedSettingsQuery).
577 WillReturnRows(pgxmock.NewRows([]string{"is_enabled", "csvlog_dest", "log_trunc", "log_dir", "lc_messages"}).
578 AddRow(true, true, false, tempDir, "en"))
579
580
581 mock.ExpectQuery(`SELECT COALESCE`).
582 WillReturnRows(pgxmock.NewRows([]string{"is_unix_socket"}).AddRow(false))
583
584
585 mock.ExpectQuery(`select name from pg_ls_logdir\(\) limit 1`).
586 WillReturnRows(pgxmock.NewRows([]string{"name"}).AddRow(logFileName))
587 mock.ExpectQuery(`select pg_read_file\(\$1, 0, 0\)`).
588 WithArgs(filepath.Join(tempDir, logFileName)).
589 WillReturnRows(pgxmock.NewRows([]string{"pg_read_file"}).AddRow(""))
590
591
592
593 mock.ExpectQuery(`select name, size, modification from pg_ls_logdir\(\) where name like '%csv' order by modification desc limit 1;`).
594 WillReturnRows(pgxmock.NewRows([]string{"name", "size", "modification"}).
595 AddRow(logFileName, int32(len(logContent)), time.Now()))
596
597 sourceConn := &sources.DbConn{
598 Source: sources.Source{
599 Name: "test-source",
600 Metrics: metrics.MetricIntervals{specialMetricServerLogEventCounts: 60},
601 },
602 Conn: mock,
603 }
604 sourceConn.RealDbname = testDbName
605
606 ctx, cancel := context.WithTimeout(testutil.TestContext, 500*time.Millisecond)
607 defer cancel()
608
609 storeCh := make(chan metrics.MeasurementEnvelope, channelBufferSize)
610
611 lp, err := NewLogParser(ctx, sourceConn, storeCh)
612 require.NoError(t, err)
613
614
615 go func() {
616 _ = lp.ParseLogs()
617 }()
618
619
620
621
622 <-ctx.Done()
623 time.Sleep(100 * time.Millisecond)
624
625
626 assert.NoError(t, mock.ExpectationsWereMet(), "All mock expectations should be met")
627
628 cancel()
629 })
630
631 t.Run("handles empty log directory gracefully", func(t *testing.T) {
632 tempDir := t.TempDir()
633 mock, err := pgxmock.NewPool()
634 require.NoError(t, err)
635 defer mock.Close()
636
637
638 mock.ExpectQuery(expectedSettingsQuery).
639 WillReturnRows(pgxmock.NewRows([]string{"is_enabled", "csvlog_dest", "log_trunc", "log_dir", "lc_messages"}).
640 AddRow(true, true, false, tempDir, "en"))
641 mock.ExpectQuery(`SELECT COALESCE`).
642 WillReturnRows(pgxmock.NewRows([]string{"is_unix_socket"}).AddRow(false))
643
644
645 mock.ExpectQuery(`select name from pg_ls_logdir\(\) limit 1`).
646 WillReturnRows(pgxmock.NewRows([]string{"name"}).AddRow(logFileName))
647 mock.ExpectQuery(`select pg_read_file\(\$1, 0, 0\)`).
648 WithArgs(filepath.Join(tempDir, logFileName)).
649 WillReturnRows(pgxmock.NewRows([]string{"pg_read_file"}).AddRow(""))
650
651
652 mock.ExpectQuery(`select name, size, modification from pg_ls_logdir\(\) where name like '%csv' order by modification desc limit 1;`).
653 WillReturnError(assert.AnError)
654
655 mock.ExpectQuery(`select name, size, modification from pg_ls_logdir\(\) where name like '%csv' order by modification desc limit 1;`).
656 WillReturnError(assert.AnError)
657
658 sourceConn := &sources.DbConn{
659 Source: sources.Source{
660 Name: "test-source",
661 Metrics: metrics.MetricIntervals{specialMetricServerLogEventCounts: 1},
662 },
663 Conn: mock,
664 }
665
666 ctx, cancel := context.WithTimeout(testutil.TestContext, 500*time.Millisecond)
667 defer cancel()
668
669 storeCh := make(chan metrics.MeasurementEnvelope, channelBufferSize)
670
671 lp, err := NewLogParser(ctx, sourceConn, storeCh)
672 require.NoError(t, err)
673
674
675 go func() {
676 _ = lp.ParseLogs()
677 }()
678
679
680 <-ctx.Done()
681 time.Sleep(100 * time.Millisecond)
682
683
684 select {
685 case m := <-storeCh:
686 t.Errorf("Expected no measurements, but received: %+v", m)
687 default:
688
689 }
690 })
691
692 t.Run("malformed CSV entries are skipped gracefully", func(t *testing.T) {
693 tempDir := t.TempDir()
694
695 malformedContent := `2023-12-01 10:30:45.123 UTC,"postgres","testdb",12345,"127.0.0.1:54321",session123,1,"SELECT",2023-12-01 10:30:00 UTC,1/234,567,ERROR,"valid entry"
696 this is not a valid CSV line at all
697 incomplete line without proper fields
698 2023-12-01 10:30:47.125 UTC,"postgres","testdb",12346,"127.0.0.1:54322",session124,1,"INSERT",2023-12-01 10:30:00 UTC,1/235,568,WARNING,"another valid entry"
699 `
700
701 mock, err := pgxmock.NewPool()
702 require.NoError(t, err)
703 defer mock.Close()
704
705
706 mock.ExpectQuery(expectedSettingsQuery).
707 WillReturnRows(pgxmock.NewRows([]string{"is_enabled", "csvlog_dest", "log_trunc", "log_dir", "lc_messages"}).
708 AddRow(true, true, false, tempDir, "en"))
709 mock.ExpectQuery(`SELECT COALESCE`).
710 WillReturnRows(pgxmock.NewRows([]string{"is_unix_socket"}).AddRow(false))
711 mock.ExpectQuery(`select name from pg_ls_logdir\(\) limit 1`).
712 WillReturnRows(pgxmock.NewRows([]string{"name"}).AddRow(logFileName))
713 mock.ExpectQuery(`select pg_read_file\(\$1, 0, 0\)`).
714 WithArgs(filepath.Join(tempDir, logFileName)).
715 WillReturnRows(pgxmock.NewRows([]string{"pg_read_file"}).AddRow(""))
716
717
718 mock.ExpectQuery(`select name, size, modification from pg_ls_logdir\(\) where name like '%csv' order by modification desc limit 1;`).
719 WillReturnRows(pgxmock.NewRows([]string{"name", "size", "modification"}).
720 AddRow(logFileName, int32(len(malformedContent)), time.Now()))
721
722 sourceConn := &sources.DbConn{
723 Source: sources.Source{
724 Name: "test-source",
725 Metrics: metrics.MetricIntervals{specialMetricServerLogEventCounts: 60},
726 },
727 Conn: mock,
728 }
729 sourceConn.RealDbname = testDbName
730
731 ctx, cancel := context.WithTimeout(testutil.TestContext, 500*time.Millisecond)
732 defer cancel()
733
734 storeCh := make(chan metrics.MeasurementEnvelope, channelBufferSize)
735
736 lp, err := NewLogParser(ctx, sourceConn, storeCh)
737 require.NoError(t, err)
738
739
740 go func() {
741 _ = lp.ParseLogs()
742 }()
743
744
745 <-ctx.Done()
746 time.Sleep(100 * time.Millisecond)
747
748
749
750
751 assert.NoError(t, mock.ExpectationsWereMet())
752
753 cancel()
754 })
755
756 t.Run("file read permission denied during parse", func(t *testing.T) {
757 tempDir := t.TempDir()
758
759 mock, err := pgxmock.NewPool()
760 require.NoError(t, err)
761 defer mock.Close()
762
763
764 mock.ExpectQuery(expectedSettingsQuery).
765 WillReturnRows(pgxmock.NewRows([]string{"is_enabled", "csvlog_dest", "log_trunc", "log_dir", "lc_messages"}).
766 AddRow(true, true, false, tempDir, "en"))
767 mock.ExpectQuery(`SELECT COALESCE`).
768 WillReturnRows(pgxmock.NewRows([]string{"is_unix_socket"}).AddRow(false))
769 mock.ExpectQuery(`select name from pg_ls_logdir\(\) limit 1`).
770 WillReturnRows(pgxmock.NewRows([]string{"name"}).AddRow(logFileName))
771 mock.ExpectQuery(`select pg_read_file\(\$1, 0, 0\)`).
772 WithArgs(filepath.Join(tempDir, logFileName)).
773 WillReturnRows(pgxmock.NewRows([]string{"pg_read_file"}).AddRow(""))
774
775
776 mock.ExpectQuery(`select name, size, modification from pg_ls_logdir\(\) where name like '%csv' order by modification desc limit 1;`).
777 WillReturnRows(pgxmock.NewRows([]string{"name", "size", "modification"}).
778 AddRow(logFileName, int32(0), time.Now()))
779
780
781 mock.ExpectQuery(`select size, modification from pg_ls_logdir\(\) where name = \$1;`).
782 WithArgs(logFileName).
783 WillReturnRows(pgxmock.NewRows([]string{"size", "modification"}).
784 AddRow(int32(len(logContent)), time.Now()))
785
786
787 mock.ExpectQuery(`select pg_read_file\(\$1, \$2, \$3\)`).
788 WithArgs(filepath.Join(tempDir, logFileName), int32(0), int32(len(logContent))).
789 WillReturnError(assert.AnError)
790
791 sourceConn := &sources.DbConn{
792 Source: sources.Source{
793 Name: "test-source",
794 Metrics: metrics.MetricIntervals{specialMetricServerLogEventCounts: 1},
795 },
796 Conn: mock,
797 }
798
799 ctx, cancel := context.WithTimeout(testutil.TestContext, 500*time.Millisecond)
800 defer cancel()
801
802 storeCh := make(chan metrics.MeasurementEnvelope, channelBufferSize)
803
804 lp, err := NewLogParser(ctx, sourceConn, storeCh)
805 require.NoError(t, err)
806
807
808 go func() {
809 _ = lp.ParseLogs()
810 }()
811
812
813 <-ctx.Done()
814 time.Sleep(100 * time.Millisecond)
815
816
817 select {
818 case m := <-storeCh:
819
820
821 data := m.Data[0]
822 assert.Equal(t, int64(0), data["error"], "Should have 0 errors since read failed")
823 assert.Equal(t, int64(0), data["warning"], "Should have 0 warnings since read failed")
824 default:
825
826 }
827 })
828 }
829
830
831
832
833
834 func TestRace_LogParserRealDbname(t *testing.T) {
835 md := sources.NewDbConn(sources.Source{Name: "race-test"})
836
837
838 lp := &LogParser{
839 LogConfig: &LogConfig{},
840 ctx: t.Context(),
841 LogsMatchRegex: regexp.MustCompile(csvLogDefaultRegEx),
842 SourceConn: md,
843 realDbname: "initial",
844 Interval: time.Second,
845 StoreCh: make(chan metrics.MeasurementEnvelope, 1),
846 eventCounts: make(map[string]int64),
847 eventCountsTotal: make(map[string]int64),
848 fileOffsets: make(map[string]uint64),
849 }
850
851 const iterations = 200
852 var wg sync.WaitGroup
853 wg.Add(2)
854
855
856 go func() {
857 defer wg.Done()
858 for range iterations {
859 md.Lock()
860 md.RealDbname = "updateddb"
861 md.Unlock()
862 }
863 }()
864
865
866 go func() {
867 defer wg.Done()
868 for range iterations {
869 _ = lp.realDbname
870 }
871 }()
872
873 wg.Wait()
874 }
875