summaryrefslogtreecommitdiff
path: root/integrationtests/dmap_test.go
diff options
context:
space:
mode:
Diffstat (limited to 'integrationtests/dmap_test.go')
-rw-r--r--integrationtests/dmap_test.go554
1 files changed, 384 insertions, 170 deletions
diff --git a/integrationtests/dmap_test.go b/integrationtests/dmap_test.go
index f772243..de57f59 100644
--- a/integrationtests/dmap_test.go
+++ b/integrationtests/dmap_test.go
@@ -15,6 +15,26 @@ func TestDMap1(t *testing.T) {
return
}
+ // Test both serverless and server modes
+ modes := []struct {
+ name string
+ useServer bool
+ }{
+ {"Serverless", false},
+ {"WithServer", true},
+ }
+
+ for _, mode := range modes {
+ t.Run(mode.name, func(t *testing.T) {
+ if err := testDMap1(t, mode.useServer); err != nil {
+ t.Error(err)
+ return
+ }
+ })
+ }
+}
+
+func testDMap1(t *testing.T, useServer bool) error {
testTable := map[string]string{
"a": "from STATS select count($line),last($time)," +
"avg($goroutines),min(concurrentConnections),max(lifetimeConnections) " +
@@ -32,68 +52,95 @@ func TestDMap1(t *testing.T) {
for subtestName, query := range testTable {
t.Log("Testing dmap with input file")
- if err := testDmap1(t, query, subtestName, false); err != nil {
+ if err := testDmap1Sub(t, query, subtestName, false, useServer); err != nil {
t.Error(err)
- return
+ return err
}
t.Log("Testing dmap with stdin input pipe")
- if err := testDmap1(t, query, subtestName, true); err != nil {
+ if err := testDmap1Sub(t, query, subtestName, true, useServer); err != nil {
t.Error(err)
- return
+ return err
}
}
+ return nil
}
-func testDmap1(t *testing.T, query, subtestName string, usePipe bool) error {
- inFile := "mapr_testdata.log"
- csvFile := fmt.Sprintf("dmap1%s.csv.tmp", subtestName)
- expectedCsvFile := fmt.Sprintf("dmap1%s.csv.expected", subtestName)
+func testDmap1Sub(t *testing.T, query, subtestName string, usePipe bool, useServer bool) error {
+ var inFile, expectedCsvFile, expectedQueryFile, csvFile string
+
+ if useServer {
+ // Use small test data for server mode to avoid channel overflow
+ inFile = "small_mapr_testdata.log"
+ csvFile = fmt.Sprintf("small_dmap1%s.csv.tmp", subtestName)
+ expectedCsvFile = fmt.Sprintf("small_dmap1%s.csv.expected", subtestName)
+ expectedQueryFile = fmt.Sprintf("small_dmap1%s.csv.query.expected", subtestName)
+ } else {
+ inFile = "mapr_testdata.log"
+ csvFile = fmt.Sprintf("dmap1%s.csv.tmp", subtestName)
+ expectedCsvFile = fmt.Sprintf("dmap1%s.csv.expected", subtestName)
+ expectedQueryFile = fmt.Sprintf("dmap1%s.csv.query.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())
- defer cancel()
-
- var stdoutCh, stderrCh <-chan string
- var cmdErrCh <-chan error
- var err error
-
- if usePipe {
- stdoutCh, stderrCh, cmdErrCh, err = startCommand(ctx, t,
- inFile, "../dmap",
- "--cfg", "none",
- "--query", query,
- "--logger", "stdout",
- "--logLevel", "info",
- "--noColor")
+ if useServer {
+ // Server mode testing
+ var args []string
+ if usePipe {
+ // For pipe mode with server, we need to handle this differently
+ // DMap with server doesn't support stdin pipe in the same way
+ // So we'll just test file mode for server
+ args = []string{"--query", query, "--logger", "stdout", "--logLevel", "info", "--noColor", inFile}
+ } else {
+ args = []string{"--query", query, "--logger", "stdout", "--logLevel", "info", "--noColor", inFile}
+ }
+ return testDMapWithServer(t, args, csvFile, expectedCsvFile, queryFile, expectedQueryFile)
} else {
- stdoutCh, stderrCh, cmdErrCh, err = startCommand(ctx, t,
- "", "../dmap",
- "--cfg", "none",
- "--query", query,
- "--logger", "stdout",
- "--logLevel", "info",
- "--noColor",
- inFile)
- }
+ // Serverless mode testing (original code)
+ ctx, cancel := context.WithCancel(context.Background())
+ defer cancel()
+
+ var stdoutCh, stderrCh <-chan string
+ var cmdErrCh <-chan error
+ var err error
+
+ if usePipe {
+ stdoutCh, stderrCh, cmdErrCh, err = startCommand(ctx, t,
+ inFile, "../dmap",
+ "--cfg", "none",
+ "--query", query,
+ "--logger", "stdout",
+ "--logLevel", "info",
+ "--noColor")
+ } else {
+ stdoutCh, stderrCh, cmdErrCh, err = startCommand(ctx, t,
+ "", "../dmap",
+ "--cfg", "none",
+ "--query", query,
+ "--logger", "stdout",
+ "--logLevel", "info",
+ "--noColor",
+ inFile)
+ }
- if err != nil {
- return err
- }
+ if err != nil {
+ return err
+ }
+
+ waitForCommand(ctx, t, stdoutCh, stderrCh, cmdErrCh)
- waitForCommand(ctx, t, stdoutCh, stderrCh, cmdErrCh)
+ if err := compareFiles(t, csvFile, expectedCsvFile); err != nil {
+ return err
+ }
+ if err := compareFiles(t, queryFile, expectedQueryFile); err != nil {
+ return err
+ }
- if err := compareFiles(t, csvFile, expectedCsvFile); err != nil {
- return err
+ os.Remove(csvFile)
+ os.Remove(queryFile)
+ return nil
}
- if err := compareFiles(t, queryFile, expectedQueryFile); err != nil {
- return err
- }
-
- os.Remove(csvFile)
- os.Remove(queryFile)
- return nil
}
func TestDMap2(t *testing.T) {
@@ -101,36 +148,73 @@ func TestDMap2(t *testing.T) {
t.Log("Skipping")
return
}
- inFile := "mapr_testdata.log"
+
+ // Test both serverless and server modes
+ modes := []struct {
+ name string
+ useServer bool
+ }{
+ {"Serverless", false},
+ {"WithServer", true},
+ }
+
+ for _, mode := range modes {
+ t.Run(mode.name, func(t *testing.T) {
+ if err := testDMap2(t, mode.useServer); err != nil {
+ t.Error(err)
+ return
+ }
+ })
+ }
+}
+
+func testDMap2(t *testing.T, useServer bool) error {
+ var inFile, expectedCsvFile, expectedQueryFile, csvFile string
outFile := "dmap2.stdout.tmp"
- csvFile := "dmap2.csv.tmp"
- expectedCsvFile := "dmap2.csv.expected"
+
+ if useServer {
+ // Use small test data for server mode to avoid channel overflow
+ inFile = "small_mapr_testdata.log"
+ csvFile = "small_dmap2.csv.tmp"
+ expectedCsvFile = "small_dmap2.csv.expected"
+ expectedQueryFile = "small_dmap2.csv.query.expected"
+ } else {
+ inFile = "mapr_testdata.log"
+ csvFile = "dmap2.csv.tmp"
+ expectedCsvFile = "dmap2.csv.expected"
+ expectedQueryFile = "dmap2.csv.query.expected"
+ }
+
queryFile := fmt.Sprintf("%s.query", csvFile)
- expectedQueryFile := "dmap2.csv.query.expected"
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)
- if err != nil {
- t.Error(err)
- return
- }
+ if useServer {
+ // Server mode testing
+ args := []string{"--query", query, "--cfg", "none", inFile}
+ return testDMapWithServer(t, args, csvFile, expectedCsvFile, queryFile, expectedQueryFile)
+ } else {
+ // Serverless mode testing (original code)
+ _, err := runCommand(context.TODO(), t, outFile,
+ "../dmap", "--query", query, "--cfg", "none", inFile)
+ if err != nil {
+ return err
+ }
- if err := compareFilesContents(t, csvFile, expectedCsvFile); err != nil {
- t.Error(err)
- return
- }
- if err := compareFiles(t, queryFile, expectedQueryFile); err != nil {
- t.Error(err)
- return
- }
+ if err := compareFilesContents(t, csvFile, expectedCsvFile); err != nil {
+ return err
+ }
+ if err := compareFiles(t, queryFile, expectedQueryFile); err != nil {
+ return err
+ }
- os.Remove(outFile)
- os.Remove(csvFile)
- os.Remove(queryFile)
+ os.Remove(outFile)
+ os.Remove(csvFile)
+ os.Remove(queryFile)
+ return nil
+ }
}
func TestDMap3(t *testing.T) {
@@ -138,56 +222,100 @@ func TestDMap3(t *testing.T) {
t.Log("Skipping")
return
}
- inFile := "mapr_testdata.log"
+
+ // Test both serverless and server modes
+ modes := []struct {
+ name string
+ useServer bool
+ }{
+ {"Serverless", false},
+ {"WithServer", true},
+ }
+
+ for _, mode := range modes {
+ t.Run(mode.name, func(t *testing.T) {
+ if err := testDMap3(t, mode.useServer); err != nil {
+ t.Error(err)
+ return
+ }
+ })
+ }
+}
+
+func testDMap3(t *testing.T, useServer bool) error {
+ var inFile, expectedCsvFile, expectedQueryFile, csvFile string
outFile := "dmap3.stdout.tmp"
- csvFile := "dmap3.csv.tmp"
- expectedCsvFile := "dmap3.csv.expected"
+
+ if useServer {
+ // Use small test data for server mode to avoid channel overflow
+ inFile = "small_mapr_testdata.log"
+ csvFile = "small_dmap3.csv.tmp"
+ expectedCsvFile = "small_dmap3.csv.expected"
+ expectedQueryFile = "small_dmap3.csv.query.expected"
+ } else {
+ inFile = "mapr_testdata.log"
+ csvFile = "dmap3.csv.tmp"
+ expectedCsvFile = "dmap3.csv.expected"
+ expectedQueryFile = "dmap3.csv.query.expected"
+ }
+
queryFile := fmt.Sprintf("%s.query", csvFile)
- expectedQueryFile := "dmap3.csv.query.expected"
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 := context.WithCancel(context.Background())
- 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)
-
- if err != nil {
- t.Error(err)
- return
- }
- waitForCommand(ctx, t, stdoutCh, stderrCh, cmdErrCh)
+ if useServer {
+ // Server mode testing - use only 3 files instead of 100 to avoid channel overflow
+ args := []string{
+ "--query", query,
+ "--cfg", "none",
+ "--logger", "stdout",
+ "--logLevel", "info",
+ "--noColor",
+ inFile, inFile, inFile,
+ }
+ return testDMapWithServer(t, args, csvFile, expectedCsvFile, queryFile, expectedQueryFile)
+ } else {
+ // Serverless mode testing (original code with 100 files)
+ ctx, cancel := context.WithCancel(context.Background())
+ defer cancel()
- if err := compareFilesContents(t, csvFile, expectedCsvFile); err != nil {
- t.Error(err)
- return
- }
- if err := compareFiles(t, queryFile, expectedQueryFile); err != nil {
- t.Error(err)
- return
- }
+ 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)
- os.Remove(outFile)
- os.Remove(csvFile)
- os.Remove(queryFile)
+ if err != nil {
+ return err
+ }
+ waitForCommand(ctx, t, stdoutCh, stderrCh, cmdErrCh)
+
+ if err := compareFilesContents(t, csvFile, expectedCsvFile); err != nil {
+ return err
+ }
+ if err := compareFiles(t, queryFile, expectedQueryFile); err != nil {
+ return err
+ }
+
+ os.Remove(outFile)
+ os.Remove(csvFile)
+ os.Remove(queryFile)
+ return nil
+ }
}
func TestDMap4Append(t *testing.T) {
@@ -195,12 +323,44 @@ func TestDMap4Append(t *testing.T) {
t.Log("Skipping")
return
}
- inFile := "mapr_testdata.log"
+
+ // Test both serverless and server modes
+ modes := []struct {
+ name string
+ useServer bool
+ }{
+ {"Serverless", false},
+ {"WithServer", true},
+ }
+
+ for _, mode := range modes {
+ t.Run(mode.name, func(t *testing.T) {
+ if err := testDMap4Append(t, mode.useServer); err != nil {
+ t.Error(err)
+ return
+ }
+ })
+ }
+}
+
+func testDMap4Append(t *testing.T, useServer bool) error {
+ var inFile, expectedCsvFile, expectedQueryFile, csvFile string
outFile := "dmap4.stdout.tmp"
- csvFile := "dmap4.csv.tmp"
- expectedCsvFile := "dmap4.csv.expected"
+
+ if useServer {
+ // Use small test data for server mode to avoid channel overflow
+ inFile = "small_mapr_testdata.log"
+ csvFile = "small_dmap4.csv.tmp"
+ expectedCsvFile = "small_dmap4.csv.expected"
+ expectedQueryFile = "small_dmap4.csv.query.expected"
+ } else {
+ inFile = "mapr_testdata.log"
+ csvFile = "dmap4.csv.tmp"
+ expectedCsvFile = "dmap4.csv.expected"
+ expectedQueryFile = "dmap4.csv.query.expected"
+ }
+
queryFile := fmt.Sprintf("%s.query", csvFile)
- expectedQueryFile := "dmap4.csv.query.expected"
// Delete in case it exists already. Otherwise, test will fail.
os.Remove(csvFile)
@@ -209,41 +369,52 @@ func TestDMap4Append(t *testing.T) {
"avg($goroutines),min($goroutines) group by $time order by count($time) "+
"outfile append %s", csvFile)
- ctx, cancel := context.WithCancel(context.Background())
- defer cancel()
-
- // 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",
+ if useServer {
+ // Server mode testing - run twice for append functionality
+ args := []string{
"--query", query,
"--cfg", "none",
"--logger", "stdout",
"--logLevel", "info",
- "--noColor", inFile)
+ "--noColor", inFile,
+ }
+ return testDMapMultipleRunsWithServer(t, args, csvFile, expectedCsvFile, queryFile, expectedQueryFile, 2)
+ } else {
+ // Serverless mode testing (original code)
+ ctx, cancel := context.WithCancel(context.Background())
+ defer cancel()
+
+ // 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)
+
+ if err != nil {
+ return err
+ }
+ waitForCommand(ctx, t, stdoutCh, stderrCh, cmdErrCh)
+ }
- if err != nil {
- t.Error(err)
- return
+ if err := compareFilesContents(t, csvFile, expectedCsvFile); err != nil {
+ return err
+ }
+ if err := compareFiles(t, queryFile, expectedQueryFile); err != nil {
+ return err
}
- 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
+ os.Remove(outFile)
+ os.Remove(csvFile)
+ os.Remove(queryFile)
+ return nil
}
-
- os.Remove(outFile)
- os.Remove(csvFile)
- os.Remove(queryFile)
}
func TestDMap5CSV(t *testing.T) {
@@ -251,12 +422,44 @@ func TestDMap5CSV(t *testing.T) {
t.Log("Skipping")
return
}
- inFile := "dmap5.csv.in"
+
+ // Test both serverless and server modes
+ modes := []struct {
+ name string
+ useServer bool
+ }{
+ {"Serverless", false},
+ {"WithServer", true},
+ }
+
+ for _, mode := range modes {
+ t.Run(mode.name, func(t *testing.T) {
+ if err := testDMap5CSV(t, mode.useServer); err != nil {
+ t.Error(err)
+ return
+ }
+ })
+ }
+}
+
+func testDMap5CSV(t *testing.T, useServer bool) error {
+ var inFile, expectedCsvFile, expectedQueryFile, csvFile string
outFile := "dmap5.stdout.tmp"
- csvFile := "dmap5.csv.tmp"
- expectedCsvFile := "dmap5.csv.expected"
+
+ if useServer {
+ // Use small test data for server mode to avoid channel overflow
+ inFile = "small_dmap5.csv.in"
+ csvFile = "small_dmap5.csv.tmp"
+ expectedCsvFile = "small_dmap5.csv.expected"
+ expectedQueryFile = "small_dmap5.csv.query.expected"
+ } else {
+ inFile = "dmap5.csv.in"
+ csvFile = "dmap5.csv.tmp"
+ expectedCsvFile = "dmap5.csv.expected"
+ expectedQueryFile = "dmap5.csv.query.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)
@@ -266,39 +469,50 @@ func TestDMap5CSV(t *testing.T) {
" set $timecount = `count($time)`, $time = `$time`, $min_goroutines = `min($goroutines)`"+
" logformat csv outfile %s", csvFile)
- ctx, cancel := context.WithCancel(context.Background())
- defer cancel()
-
- // 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",
+ if useServer {
+ // Server mode testing - run twice (CSV input format with append)
+ args := []string{
"--query", query,
"--cfg", "none",
"--logger", "stdout",
"--logLevel", "info",
- "--noColor", inFile)
+ "--noColor", inFile,
+ }
+ return testDMapMultipleRunsWithServer(t, args, csvFile, expectedCsvFile, queryFile, expectedQueryFile, 2)
+ } else {
+ // Serverless mode testing (original code)
+ ctx, cancel := context.WithCancel(context.Background())
+ defer cancel()
+
+ // 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)
+
+ if err != nil {
+ return err
+ }
+ waitForCommand(ctx, t, stdoutCh, stderrCh, cmdErrCh)
+ }
- if err != nil {
- t.Error(err)
- return
+ if err := compareFilesContents(t, csvFile, expectedCsvFile); err != nil {
+ return err
+ }
+ if err := compareFiles(t, queryFile, expectedQueryFile); err != nil {
+ return err
}
- waitForCommand(ctx, t, stdoutCh, stderrCh, cmdErrCh)
- }
- if err := compareFilesContents(t, csvFile, expectedCsvFile); err != nil {
- t.Error(err)
- return
+ os.Remove(outFile)
+ os.Remove(csvFile)
+ os.Remove(queryFile)
+ return nil
}
- if err := compareFiles(t, queryFile, expectedQueryFile); err != nil {
- t.Error(err)
- return
- }
-
- os.Remove(outFile)
- os.Remove(csvFile)
- os.Remove(queryFile)
}