diff options
| author | Paul Buetow <paul@buetow.org> | 2026-07-22 23:51:28 +0300 |
|---|---|---|
| committer | Paul Buetow <paul@buetow.org> | 2026-07-22 23:51:28 +0300 |
| commit | bc2767c87c4090798c4c7d15e101ed066e947301 (patch) | |
| tree | 0f47a2e0b159a617eb56509bbca1c998196453c4 /integrationtests/dmap_test.go | |
| parent | 849951be1d1a7ee9f9302006ccb187bf5b4e36f3 (diff) | |
test: DTail fork — integration test suite and fixtures
Squashed development of the integration test suite (integrationtests/) covering
DCat, DGrep, DMap (serverless + server mode), DTail follow, DServer, DTailHealth,
journal source reads, auth-key fast reconnect, interactive query reload,
client-deadline/timeout behaviour, and the single-mode read/output path.
Includes real test fixtures (dserver*.cfg, dmap_csv_multifile_*.csv.in,
test_server_*.json, *.expected golden files) and deterministic synchronization
helpers (waitContains-style barriers) replacing racy fixed-timing assertions.
Accidental debug/output dumps that earlier commits added here (captured client
output, strace logs, ad-hoc turbo_test_output/manual_output/test_output files,
throwaway debug scripts) are intentionally excluded and gitignored.
Co-Authored-By: Claude Opus 4.8 <noreply@anthropic.com>
Diffstat (limited to 'integrationtests/dmap_test.go')
| -rw-r--r-- | integrationtests/dmap_test.go | 754 |
1 files changed, 589 insertions, 165 deletions
diff --git a/integrationtests/dmap_test.go b/integrationtests/dmap_test.go index f772243..c9f4ecb 100644 --- a/integrationtests/dmap_test.go +++ b/integrationtests/dmap_test.go @@ -3,17 +3,15 @@ package integrationtests import ( "context" "fmt" - "os" + "path/filepath" "testing" - - "github.com/mimecast/dtail/internal/config" ) func TestDMap1(t *testing.T) { - if !config.Env("DTAIL_INTEGRATION_TEST_RUN_MODE") { - t.Log("Skipping") - return - } + skipIfNotIntegrationTest(t) + cleanupTmpFiles(t) + testLogger := NewTestLogger("TestDMap1") + defer testLogger.WriteLogFile() testTable := map[string]string{ "a": "from STATS select count($line),last($time)," + @@ -30,275 +28,701 @@ func TestDMap1(t *testing.T) { "$foo = 42, $bar = \"baz\", $baz = $time group by $hostname", } - for subtestName, query := range testTable { - t.Log("Testing dmap with input file") - if err := testDmap1(t, query, subtestName, false); err != nil { - t.Error(err) - return + // Test in serverless mode + t.Run("Serverless", func(t *testing.T) { + for subtestName, query := range testTable { + t.Run(subtestName, func(t *testing.T) { + t.Log("Testing dmap with input file") + testDmap1Serverless(t, testLogger, query, subtestName, false) + + t.Log("Testing dmap with stdin input pipe") + testDmap1Serverless(t, testLogger, query, subtestName, true) + }) } - t.Log("Testing dmap with stdin input pipe") - if err := testDmap1(t, query, subtestName, true); err != nil { - t.Error(err) - return + }) + + // Test in server mode + t.Run("ServerMode", func(t *testing.T) { + for subtestName, query := range testTable { + t.Run(subtestName, func(t *testing.T) { + t.Log("Testing dmap with input file in server mode") + testDmap1WithServer(t, testLogger, query, subtestName) + }) } - } + }) } -func testDmap1(t *testing.T, query, subtestName string, usePipe bool) error { - inFile := "mapr_testdata.log" +func testDmap1Serverless(t *testing.T, logger *TestLogger, query, subtestName string, usePipe bool) { + paths := GetStandardTestPaths() csvFile := fmt.Sprintf("dmap1%s.csv.tmp", subtestName) expectedCsvFile := fmt.Sprintf("dmap1%s.csv.expected", subtestName) queryFile := fmt.Sprintf("%s.query", csvFile) - expectedQueryFile := fmt.Sprintf("dmap1%s.csv.query.expected", subtestName) query = fmt.Sprintf("%s outfile %s", query, csvFile) - ctx, cancel := context.WithCancel(context.Background()) + cleanupFiles(t, csvFile, queryFile) + + ctxTimeout, cancel := createTestContextWithTimeout(t) + ctx := WithTestLogger(ctxTimeout, logger) defer cancel() var stdoutCh, stderrCh <-chan string var cmdErrCh <-chan error var err error + args := NewCommandArgs() + args.Logger = "stdout" + args.LogLevel = "info" + args.NoColor = true + args.ExtraArgs = []string{"--query", query} + if usePipe { stdoutCh, stderrCh, cmdErrCh, err = startCommand(ctx, t, - inFile, "../dmap", - "--cfg", "none", - "--query", query, - "--logger", "stdout", - "--logLevel", "info", - "--noColor") + paths.MaprTestData, "../dmap", args.ToSlice()...) } else { stdoutCh, stderrCh, cmdErrCh, err = startCommand(ctx, t, - "", "../dmap", - "--cfg", "none", - "--query", query, - "--logger", "stdout", - "--logLevel", "info", - "--noColor", - inFile) + "", "../dmap", append(args.ToSlice(), paths.MaprTestData)...) } if err != nil { - return err + t.Error(err) + return } waitForCommand(ctx, t, stdoutCh, stderrCh, cmdErrCh) - if err := compareFiles(t, csvFile, expectedCsvFile); err != nil { - return err + if err := compareFilesContentsWithContext(ctx, t, csvFile, expectedCsvFile); err != nil { + t.Error(err) + } + if err := verifyQueryFile(t, queryFile, query); err != nil { + t.Error(err) } - if err := compareFiles(t, queryFile, expectedQueryFile); err != nil { - return err +} + +func testDmap1WithServer(t *testing.T, logger *TestLogger, query, subtestName string) { + ctx := WithTestLogger(context.Background(), logger) + paths := GetStandardTestPaths() + csvFile := fmt.Sprintf("dmap1%s.csv.tmp", subtestName) + expectedCsvFile := fmt.Sprintf("dmap1%s.csv.expected", subtestName) + queryFile := fmt.Sprintf("%s.query", csvFile) + query = fmt.Sprintf("%s outfile %s", query, csvFile) + + cleanupFiles(t, csvFile, queryFile) + + server := NewTestServer(t) + if err := server.Start("error"); err != nil { + t.Error(err) + return + } + + args := NewCommandArgs() + args.Logger = "stdout" + args.LogLevel = "info" + args.NoColor = true + args.Servers = []string{server.Address()} + args.TrustAllHosts = true + args.Files = []string{paths.MaprTestData} + args.ExtraArgs = []string{"--query", query} + + stdoutCh, stderrCh, cmdErrCh, err := startCommand(server.ctx, t, + "", "../dmap", args.ToSlice()...) + if err != nil { + t.Error(err) + return } - os.Remove(csvFile) - os.Remove(queryFile) - return nil + waitForCommand(server.ctx, t, stdoutCh, stderrCh, cmdErrCh) + + if err := compareFilesContentsWithContext(ctx, t, csvFile, expectedCsvFile); err != nil { + t.Error(err) + } + if err := verifyQueryFile(t, queryFile, query); err != nil { + t.Error(err) + } } func TestDMap2(t *testing.T) { - if !config.Env("DTAIL_INTEGRATION_TEST_RUN_MODE") { - t.Log("Skipping") + cleanupTmpFiles(t) + testLogger := NewTestLogger("TestDMap2") + defer testLogger.WriteLogFile() + runDualModeTest(t, DualModeTest{ + Name: "TestDMap2", + ServerlessTest: func(t *testing.T) { testDMap2Serverless(t, testLogger) }, + ServerTest: func(t *testing.T) { testDMap2WithServer(t, testLogger) }, + }) +} + +func TestDMapOutfileArbitraryPath(t *testing.T) { + cleanupTmpFiles(t) + testLogger := NewTestLogger("TestDMapOutfileArbitraryPath") + defer testLogger.WriteLogFile() + runDualModeTest(t, DualModeTest{ + Name: "TestDMapOutfileArbitraryPath", + ServerlessTest: func(t *testing.T) { + testDMapOutfileArbitraryPathServerless(t, testLogger) + }, + ServerTest: func(t *testing.T) { + testDMapOutfileArbitraryPathWithServer(t, testLogger) + }, + }) +} + +func testDMapOutfileArbitraryPathServerless(t *testing.T, logger *TestLogger) { + paths := GetStandardTestPaths() + outFile := "dmap_outfile_arbitrary_serverless.stdout.tmp" + csvFile := filepath.Join(t.TempDir(), "dmap_outfile_arbitrary_serverless.csv") + expectedCsvFile := "dmap2.csv.expected" + queryFile := fmt.Sprintf("%s.query", csvFile) + query := fmt.Sprintf("from STATS select count($time),$time,max($goroutines),"+ + "avg($goroutines),min($goroutines) group by $time order by count($time) "+ + "outfile %s", csvFile) + + if !filepath.IsAbs(csvFile) { + t.Fatalf("expected absolute outfile path, got %q", csvFile) + } + cleanupFiles(t, outFile, csvFile, queryFile) + + ctxTimeout, cancel := createTestContextWithTimeout(t) + ctx := WithTestLogger(ctxTimeout, logger) + defer cancel() + + _, err := runCommand(ctx, t, outFile, + "../dmap", "--query", query, "--cfg", "none", paths.MaprTestData) + if err != nil { + t.Error(err) + return + } + + if err := compareFilesContentsWithContext(ctx, t, csvFile, expectedCsvFile); err != nil { + t.Error(err) + } + if err := verifyQueryFile(t, queryFile, query); err != nil { + t.Error(err) + } +} + +func testDMapOutfileArbitraryPathWithServer(t *testing.T, logger *TestLogger) { + ctx := WithTestLogger(context.Background(), logger) + paths := GetStandardTestPaths() + outFile := "dmap_outfile_arbitrary_server.stdout.tmp" + csvFile := filepath.Join(t.TempDir(), "dmap_outfile_arbitrary_server.csv") + expectedCsvFile := "dmap2.csv.expected" + queryFile := fmt.Sprintf("%s.query", csvFile) + query := fmt.Sprintf("from STATS select count($time),$time,max($goroutines),"+ + "avg($goroutines),min($goroutines) group by $time order by count($time) "+ + "outfile %s", csvFile) + + if !filepath.IsAbs(csvFile) { + t.Fatalf("expected absolute outfile path, got %q", csvFile) + } + cleanupFiles(t, outFile, csvFile, queryFile) + + server := NewTestServer(t) + if err := server.Start("error"); err != nil { + t.Error(err) return } - inFile := "mapr_testdata.log" - outFile := "dmap2.stdout.tmp" - csvFile := "dmap2.csv.tmp" + + args := NewCommandArgs() + args.Servers = []string{server.Address()} + args.TrustAllHosts = true + args.NoColor = true + args.Files = []string{paths.MaprTestData} + args.ExtraArgs = []string{"--query", query} + + _, err := runCommand(server.ctx, t, outFile, + "../dmap", args.ToSlice()...) + if err != nil { + t.Error(err) + return + } + + if err := compareFilesContentsWithContext(ctx, t, csvFile, expectedCsvFile); err != nil { + t.Error(err) + } + if err := verifyQueryFile(t, queryFile, query); err != nil { + t.Error(err) + } +} + +func testDMap2Serverless(t *testing.T, logger *TestLogger) { + paths := GetStandardTestPaths() + outFile := "dmap2_serverless.stdout.tmp" + csvFile := "dmap2_serverless.csv.tmp" expectedCsvFile := "dmap2.csv.expected" queryFile := fmt.Sprintf("%s.query", csvFile) - expectedQueryFile := "dmap2.csv.query.expected" + cleanupFiles(t, outFile, csvFile, queryFile) query := fmt.Sprintf("from STATS select count($time),$time,max($goroutines),"+ "avg($goroutines),min($goroutines) group by $time order by count($time) "+ "outfile %s", csvFile) - _, err := runCommand(context.TODO(), t, outFile, - "../dmap", "--query", query, "--cfg", "none", inFile) + ctxTimeout, cancel := createTestContextWithTimeout(t) + ctx := WithTestLogger(ctxTimeout, logger) + defer cancel() + _, err := runCommand(ctx, t, outFile, + "../dmap", "--query", query, "--cfg", "none", paths.MaprTestData) if err != nil { t.Error(err) return } - if err := compareFilesContents(t, csvFile, expectedCsvFile); err != nil { + if err := compareFilesContentsWithContext(ctx, t, csvFile, expectedCsvFile); err != nil { + t.Error(err) + } + if err := verifyQueryFile(t, queryFile, query); err != nil { + t.Error(err) + } +} + +func testDMap2WithServer(t *testing.T, logger *TestLogger) { + ctx := WithTestLogger(context.Background(), logger) + paths := GetStandardTestPaths() + outFile := "dmap2_server.stdout.tmp" + csvFile := "dmap2_server.csv.tmp" + expectedCsvFile := "dmap2.csv.expected" + queryFile := fmt.Sprintf("%s.query", csvFile) + cleanupFiles(t, outFile, csvFile, queryFile) + + server := NewTestServer(t) + if err := server.Start("error"); err != nil { t.Error(err) return } - if err := compareFiles(t, queryFile, expectedQueryFile); err != nil { + + query := fmt.Sprintf("from STATS select count($time),$time,max($goroutines),"+ + "avg($goroutines),min($goroutines) group by $time order by count($time) "+ + "outfile %s", csvFile) + + args := NewCommandArgs() + args.Servers = []string{server.Address()} + args.TrustAllHosts = true + args.NoColor = true + args.Files = []string{paths.MaprTestData} + args.ExtraArgs = []string{"--query", query} + + _, err := runCommand(server.ctx, t, outFile, + "../dmap", args.ToSlice()...) + if err != nil { t.Error(err) return } - os.Remove(outFile) - os.Remove(csvFile) - os.Remove(queryFile) + if err := compareFilesContentsWithContext(ctx, t, csvFile, expectedCsvFile); err != nil { + t.Error(err) + } + if err := verifyQueryFile(t, queryFile, query); err != nil { + t.Error(err) + } } func TestDMap3(t *testing.T) { - if !config.Env("DTAIL_INTEGRATION_TEST_RUN_MODE") { - t.Log("Skipping") - return - } - inFile := "mapr_testdata.log" - outFile := "dmap3.stdout.tmp" - csvFile := "dmap3.csv.tmp" + cleanupTmpFiles(t) + testLogger := NewTestLogger("TestDMap3") + defer testLogger.WriteLogFile() + runDualModeTest(t, DualModeTest{ + Name: "TestDMap3", + ServerlessTest: func(t *testing.T) { testDMap3Serverless(t, testLogger) }, + ServerTest: func(t *testing.T) { testDMap3WithServer(t, testLogger) }, + }) +} + +func testDMap3Serverless(t *testing.T, logger *TestLogger) { + paths := GetStandardTestPaths() + outFile := "dmap3_serverless.stdout.tmp" + csvFile := "dmap3_serverless.csv.tmp" expectedCsvFile := "dmap3.csv.expected" queryFile := fmt.Sprintf("%s.query", csvFile) - expectedQueryFile := "dmap3.csv.query.expected" + cleanupFiles(t, outFile, csvFile, queryFile) - query := fmt.Sprintf("from STATS select count($time),$time,max($goroutines),"+ - "avg($goroutines),min($goroutines) group by $time order by count($time) "+ + query := fmt.Sprintf("from STATS select count($time),$time,max($goroutines),avg($goroutines),min($goroutines) "+ + "group by $time order by count($time) desc "+ "outfile %s", csvFile) - ctx, cancel := context.WithCancel(context.Background()) + // Create a large list of input files + var inputFiles []string + for i := 0; i < 100; i++ { + inputFiles = append(inputFiles, paths.MaprTestData) + } + + // Simply run dmap with multiple input files directly + // Use longer timeout for processing 100 files + ctxTimeout, cancel := createTestContextWithLongTimeout(t) + ctx := WithTestLogger(ctxTimeout, logger) defer cancel() - stdoutCh, stderrCh, cmdErrCh, err := startCommand(ctx, t, - "", "../dmap", - "--query", query, - "--cfg", "none", - "--logger", "stdout", - "--logLevel", "info", - "--noColor", - inFile, inFile, inFile, inFile, inFile, inFile, inFile, inFile, inFile, inFile, - inFile, inFile, inFile, inFile, inFile, inFile, inFile, inFile, inFile, inFile, - inFile, inFile, inFile, inFile, inFile, inFile, inFile, inFile, inFile, inFile, - inFile, inFile, inFile, inFile, inFile, inFile, inFile, inFile, inFile, inFile, - inFile, inFile, inFile, inFile, inFile, inFile, inFile, inFile, inFile, inFile, - inFile, inFile, inFile, inFile, inFile, inFile, inFile, inFile, inFile, inFile, - inFile, inFile, inFile, inFile, inFile, inFile, inFile, inFile, inFile, inFile, - inFile, inFile, inFile, inFile, inFile, inFile, inFile, inFile, inFile, inFile, - inFile, inFile, inFile, inFile, inFile, inFile, inFile, inFile, inFile, inFile, - inFile, inFile, inFile, inFile, inFile, inFile, inFile, inFile, inFile, inFile) + args := NewCommandArgs() + args.ExtraArgs = []string{"--query", query} + _, err := runCommand(ctx, t, outFile, + "../dmap", append(args.ToSlice(), inputFiles...)...) if err != nil { t.Error(err) return } - waitForCommand(ctx, t, stdoutCh, stderrCh, cmdErrCh) - if err := compareFilesContents(t, csvFile, expectedCsvFile); err != nil { + if err := compareFilesContentsWithContext(ctx, t, csvFile, expectedCsvFile); err != nil { + t.Error(err) + } + if err := verifyQueryFile(t, queryFile, query); err != nil { + t.Error(err) + } +} + +func testDMap3WithServer(t *testing.T, logger *TestLogger) { + ctx := WithTestLogger(context.Background(), logger) + paths := GetStandardTestPaths() + outFile := "dmap3_server.stdout.tmp" + csvFile := "dmap3_server.csv.tmp" + expectedCsvFile := "dmap3.csv.expected" + queryFile := fmt.Sprintf("%s.query", csvFile) + cleanupFiles(t, outFile, csvFile, queryFile) + + server := NewTestServer(t) + // Correctness test for server-mode MapReduce over 100 files: the CSV output + // must equal the expected file. This is not a turbo-vs-non-turbo comparison. + // Under DTAIL_INTEGRATION_TEST_RUN_MODE the server force-disables turbo boost + // (see internal/config/initializer.go), so this dserver always runs the + // non-turbo aggregate path regardless of any DTAIL_TURBOBOOST_* env var. + // Turbo aggregate correctness is covered by the unit tests in + // internal/mapr/server/turbo_aggregate_test.go. + cfg := &ServerConfig{ + Port: server.port, + BindAddress: server.bindAddress, + LogLevel: "error", + ExtraArgs: []string{"--cfg", "test_server_100files.json"}, + } + if err := server.StartWithConfig(cfg); err != nil { t.Error(err) return } - if err := compareFiles(t, queryFile, expectedQueryFile); err != nil { + + query := fmt.Sprintf("from STATS select count($time),$time,max($goroutines),avg($goroutines),min($goroutines) "+ + "group by $time order by count($time) desc "+ + "outfile %s", csvFile) + + // Create a large list of input files + var inputFiles []string + for i := 0; i < 100; i++ { + inputFiles = append(inputFiles, paths.MaprTestData) + } + + args := NewCommandArgs() + args.Servers = []string{server.Address()} + args.TrustAllHosts = true + args.NoColor = true + args.Files = inputFiles + args.ExtraArgs = []string{"--query", query} + + _, err := runCommand(server.ctx, t, outFile, + "../dmap", args.ToSlice()...) + if err != nil { t.Error(err) return } - os.Remove(outFile) - os.Remove(csvFile) - os.Remove(queryFile) + if err := compareFilesContentsWithContext(ctx, t, csvFile, expectedCsvFile); err != nil { + t.Error(err) + } + if err := verifyQueryFile(t, queryFile, query); err != nil { + t.Error(err) + } } func TestDMap4Append(t *testing.T) { - if !config.Env("DTAIL_INTEGRATION_TEST_RUN_MODE") { - t.Log("Skipping") - return - } - inFile := "mapr_testdata.log" - outFile := "dmap4.stdout.tmp" - csvFile := "dmap4.csv.tmp" - expectedCsvFile := "dmap4.csv.expected" + cleanupTmpFiles(t) + testLogger := NewTestLogger("TestDMap4Append") + defer testLogger.WriteLogFile() + runDualModeTest(t, DualModeTest{ + Name: "TestDMap4Append", + ServerlessTest: func(t *testing.T) { testDMap4AppendServerless(t, testLogger) }, + ServerTest: func(t *testing.T) { testDMap4AppendWithServer(t, testLogger) }, + }) +} + +func testDMap4AppendServerless(t *testing.T, logger *TestLogger) { + paths := GetStandardTestPaths() + csvFile := "dmap4_serverless.csv.tmp" queryFile := fmt.Sprintf("%s.query", csvFile) - expectedQueryFile := "dmap4.csv.query.expected" - // Delete in case it exists already. Otherwise, test will fail. - os.Remove(csvFile) + // Clean up files once at the beginning + cleanupFiles(t, csvFile, queryFile) - query := fmt.Sprintf("from STATS select count($time),$time,max($goroutines),"+ - "avg($goroutines),min($goroutines) group by $time order by count($time) "+ - "outfile append %s", csvFile) + t.Run("FirstQuery", func(t *testing.T) { + stdout := "dmap4_serverless.stdout1.tmp" + cleanupFiles(t, stdout) - ctx, cancel := context.WithCancel(context.Background()) - defer cancel() + // First query + query := fmt.Sprintf("from STATS select count($time),$time,max($goroutines),"+ + "avg($goroutines),min($goroutines) group by $time order by count($time) "+ + "outfile %s", csvFile) + + ctx, cancel := createTestContextWithTimeout(t) + defer cancel() + _, err := runCommand(ctx, t, stdout, + "../dmap", "--query", query, "--cfg", "none", paths.MaprTestData) + if err != nil { + t.Error(err) + return + } - // Run dmap command twice, it should append in the 2nd iteration the new results to the already existing - // file as we specified "outfile append". That works transparently for any mapreduce query - // (e.g. also for the dtail command in streaming mode). But it is easier to test with the dmap - // command. - for i := 0; i < 2; i++ { - stdoutCh, stderrCh, cmdErrCh, err := startCommand(ctx, t, - "", "../dmap", - "--query", query, - "--cfg", "none", - "--logger", "stdout", - "--logLevel", "info", - "--noColor", inFile) + // Verify the CSV output + if err := compareFilesContentsWithContext(ctx, t, csvFile, "dmap4_query1.csv.expected"); err != nil { + t.Error(err) + } + // Verify the query file + if err := verifyQueryFile(t, queryFile, query); err != nil { + t.Error(err) + } + }) + + t.Run("SecondQueryWithAppend", func(t *testing.T) { + stdout := "dmap4_serverless.stdout2.tmp" + cleanupFiles(t, stdout) + + // Second query with append + query := fmt.Sprintf("from STATS select count($time),$time,max($goroutines),"+ + "avg($goroutines),min($goroutines) group by $time order by avg($goroutines) reverse "+ + "outfile append:%s", csvFile) + + ctx, cancel := createTestContextWithTimeout(t) + defer cancel() + _, err := runCommand(ctx, t, stdout, + "../dmap", "--query", query, "--cfg", "none", paths.MaprTestData) if err != nil { t.Error(err) return } - waitForCommand(ctx, t, stdoutCh, stderrCh, cmdErrCh) - } - if err := compareFilesContents(t, csvFile, expectedCsvFile); err != nil { - t.Error(err) - return - } - if err := compareFiles(t, queryFile, expectedQueryFile); err != nil { - t.Error(err) - return - } + // Verify the CSV output (should still be the first query result - append doesn't change existing file) + if err := compareFilesContentsWithContext(ctx, t, csvFile, "dmap4_query1.csv.expected"); err != nil { + t.Error(err) + } + }) + + t.Run("ThirdQueryWithAppend", func(t *testing.T) { + stdout := "dmap4_serverless.stdout3.tmp" + cleanupFiles(t, stdout) - os.Remove(outFile) - os.Remove(csvFile) - os.Remove(queryFile) + // Third query with append (different structure) + query := fmt.Sprintf("from STATS select count($line),$hostname "+ + "group by $hostname "+ + "outfile append:%s", csvFile) + + ctx, cancel := createTestContextWithTimeout(t) + defer cancel() + _, err := runCommand(ctx, t, stdout, + "../dmap", "--query", query, "--cfg", "none", paths.MaprTestData) + if err != nil { + t.Error(err) + return + } + + // Verify the CSV output (should still be the first query result - append doesn't change existing file) + if err := compareFilesContentsWithContext(ctx, t, csvFile, "dmap4_query1.csv.expected"); err != nil { + t.Error(err) + } + + // For append test, the query file should still contain the first query + firstQuery := fmt.Sprintf("from STATS select count($time),$time,max($goroutines),"+ + "avg($goroutines),min($goroutines) group by $time order by count($time) "+ + "outfile %s", csvFile) + if err := verifyQueryFile(t, queryFile, firstQuery); err != nil { + t.Error(err) + } + }) } -func TestDMap5CSV(t *testing.T) { - if !config.Env("DTAIL_INTEGRATION_TEST_RUN_MODE") { - t.Log("Skipping") +func testDMap4AppendWithServer(t *testing.T, logger *TestLogger) { + ctx := WithTestLogger(context.Background(), logger) + paths := GetStandardTestPaths() + csvFile := "dmap4_server.csv.tmp" + queryFile := fmt.Sprintf("%s.query", csvFile) + + server := NewTestServer(t) + if err := server.Start("error"); err != nil { + t.Error(err) return } - inFile := "dmap5.csv.in" - outFile := "dmap5.stdout.tmp" - csvFile := "dmap5.csv.tmp" - expectedCsvFile := "dmap5.csv.expected" - queryFile := fmt.Sprintf("%s.query", csvFile) - expectedQueryFile := "dmap5.csv.query.expected" - // Delete in case it exists already. Otherwise, test will fail. - os.Remove(csvFile) + baseArgs := NewCommandArgs() + baseArgs.Servers = []string{server.Address()} + baseArgs.TrustAllHosts = true + baseArgs.NoColor = true + baseArgs.Files = []string{paths.MaprTestData} - query := fmt.Sprintf("select sum($timecount),last($time),min($min_goroutines),"+ - " group by $hostname"+ - " set $timecount = `count($time)`, $time = `$time`, $min_goroutines = `min($goroutines)`"+ - " logformat csv outfile %s", csvFile) + // Clean up files once at the beginning + cleanupFiles(t, csvFile, queryFile) - ctx, cancel := context.WithCancel(context.Background()) - defer cancel() + t.Run("FirstQuery", func(t *testing.T) { + stdout := "dmap4_server.stdout1.tmp" + cleanupFiles(t, stdout) + + // First query + query := fmt.Sprintf("from STATS select count($time),$time,max($goroutines),"+ + "avg($goroutines),min($goroutines) group by $time order by count($time) "+ + "outfile %s", csvFile) - // Run dmap command twice, it should append in the 2nd iteration the new results to the already existing - // file as we specified "outfile append". That works transparently for any mapreduce query - // (e.g. also for the dtail command in streaming mode). But it is easier to test with the dmap - // command. - for i := 0; i < 2; i++ { - stdoutCh, stderrCh, cmdErrCh, err := startCommand(ctx, t, - "", "../dmap", - "--query", query, - "--cfg", "none", - "--logger", "stdout", - "--logLevel", "info", - "--noColor", inFile) + args := *baseArgs + args.ExtraArgs = []string{"--query", query} + _, err := runCommand(server.ctx, t, stdout, + "../dmap", args.ToSlice()...) if err != nil { t.Error(err) return } - waitForCommand(ctx, t, stdoutCh, stderrCh, cmdErrCh) + + // Verify the CSV output + if err := compareFilesContentsWithContext(ctx, t, csvFile, "dmap4_query1.csv.expected"); err != nil { + t.Error(err) + } + + // Verify the query file + if err := verifyQueryFile(t, queryFile, query); err != nil { + t.Error(err) + } + }) + + t.Run("SecondQueryWithAppend", func(t *testing.T) { + stdout := "dmap4_server.stdout2.tmp" + cleanupFiles(t, stdout) + + // Second query with append + query := fmt.Sprintf("from STATS select count($time),$time,max($goroutines),"+ + "avg($goroutines),min($goroutines) group by $time order by avg($goroutines) reverse "+ + "outfile append:%s", csvFile) + + args := *baseArgs + args.ExtraArgs = []string{"--query", query} + + _, err := runCommand(server.ctx, t, stdout, + "../dmap", args.ToSlice()...) + if err != nil { + t.Error(err) + return + } + + // Verify the CSV output (should still be the first query result - append doesn't change existing file) + if err := compareFilesContentsWithContext(ctx, t, csvFile, "dmap4_query1.csv.expected"); err != nil { + t.Error(err) + } + }) + + t.Run("ThirdQueryWithAppend", func(t *testing.T) { + stdout := "dmap4_server.stdout3.tmp" + cleanupFiles(t, stdout) + + // Third query with append (different structure) + query := fmt.Sprintf("from STATS select count($line),$hostname "+ + "group by $hostname "+ + "outfile append:%s", csvFile) + + args := *baseArgs + args.ExtraArgs = []string{"--query", query} + + _, err := runCommand(server.ctx, t, stdout, + "../dmap", args.ToSlice()...) + if err != nil { + t.Error(err) + return + } + + // Verify the CSV output (should still be the first query result - append doesn't change existing file) + if err := compareFilesContentsWithContext(ctx, t, csvFile, "dmap4_query1.csv.expected"); err != nil { + t.Error(err) + } + + // For append test, the query file should still contain the first query + firstQuery := fmt.Sprintf("from STATS select count($time),$time,max($goroutines),"+ + "avg($goroutines),min($goroutines) group by $time order by count($time) "+ + "outfile %s", csvFile) + if err := verifyQueryFile(t, queryFile, firstQuery); err != nil { + t.Error(err) + } + }) +} + +func TestDMap5CSV(t *testing.T) { + cleanupTmpFiles(t) + testLogger := NewTestLogger("TestDMap5CSV") + defer testLogger.WriteLogFile() + runDualModeTest(t, DualModeTest{ + Name: "TestDMap5CSV", + ServerlessTest: func(t *testing.T) { testDMap5CSVServerless(t, testLogger) }, + ServerTest: func(t *testing.T) { testDMap5CSVWithServer(t, testLogger) }, + }) +} + +func testDMap5CSVServerless(t *testing.T, logger *TestLogger) { + inFile := "dmap5.csv.in" + csvFile := "dmap5_serverless.csv.tmp" + expectedCsvFile := "dmap5.csv.expected" + queryFile := fmt.Sprintf("%s.query", csvFile) + outFile := "dmap5_serverless.stdout.tmp" + cleanupFiles(t, csvFile, queryFile, outFile) + + query := fmt.Sprintf("select sum($timecount),last($time),min($min_goroutines) "+ + "group by $hostname set $timecount = `count($time)`, $time = `$time`, "+ + "$min_goroutines = `min($goroutines)` logformat csv outfile %s", csvFile) + + ctxTimeout, cancel := createTestContextWithTimeout(t) + ctx := WithTestLogger(ctxTimeout, logger) + defer cancel() + _, err := runCommand(ctx, t, outFile, + "../dmap", "--query", query, "--cfg", "none", inFile) + if err != nil { + t.Error(err) + return + } + + if err := compareFilesContentsWithContext(ctx, t, csvFile, expectedCsvFile); err != nil { + t.Error(err) } + // Verify the query file contains the expected query + if err := verifyQueryFile(t, queryFile, query); err != nil { + t.Error(err) + } +} + +func testDMap5CSVWithServer(t *testing.T, logger *TestLogger) { + ctx := WithTestLogger(context.Background(), logger) + inFile := "dmap5.csv.in" + csvFile := "dmap5_server.csv.tmp" + expectedCsvFile := "dmap5.csv.expected" + queryFile := fmt.Sprintf("%s.query", csvFile) + outFile := "dmap5_server.stdout.tmp" + cleanupFiles(t, csvFile, queryFile, outFile) - if err := compareFilesContents(t, csvFile, expectedCsvFile); err != nil { + server := NewTestServer(t) + if err := server.Start("error"); err != nil { t.Error(err) return } - if err := compareFiles(t, queryFile, expectedQueryFile); err != nil { + + query := fmt.Sprintf("select sum($timecount),last($time),min($min_goroutines) "+ + "group by $hostname set $timecount = `count($time)`, $time = `$time`, "+ + "$min_goroutines = `min($goroutines)` logformat csv outfile %s", csvFile) + + args := NewCommandArgs() + args.Servers = []string{server.Address()} + args.TrustAllHosts = true + args.NoColor = true + args.Files = []string{inFile} + args.ExtraArgs = []string{"--query", query} + + _, err := runCommand(server.ctx, t, outFile, + "../dmap", args.ToSlice()...) + if err != nil { t.Error(err) return } - os.Remove(outFile) - os.Remove(csvFile) - os.Remove(queryFile) + if err := compareFilesContentsWithContext(ctx, t, csvFile, expectedCsvFile); err != nil { + t.Error(err) + } + // Verify the query file contains the expected query + if err := verifyQueryFile(t, queryFile, query); err != nil { + t.Error(err) + } } |
