summaryrefslogtreecommitdiff
path: root/internal/server/handlers/mapcommand_race_test.go
blob: 02c3078a40f492fcc9b8f2e1645be4fd1d54a964 (plain)
1
2
3
4
5
6
7
8
9
10
11
12
13
14
15
16
17
18
19
20
21
22
23
24
25
26
27
28
29
30
31
32
33
34
35
36
37
38
39
40
41
42
43
44
45
46
47
48
49
50
51
52
53
54
55
56
57
58
59
60
61
62
63
64
package handlers

// TestAggregatePointerRaceWithShutdown is a negative test that exercises
// the race between handleMapCommand writing h.aggregate and baseHandler
// Shutdown reading that same pointer concurrently. Running this test with -race
// detects unsynchronized access before the atomic.Pointer fix is applied, and
// passes cleanly after. The regular-aggregate counterpart was removed with the
// regular server.Aggregate itself (task hv0); output is now the only aggregate.
//
// The test does not set up a real MapReduce query; instead it injects a real
// output aggregate via the atomic accessors so the test is self-contained.

import (
	"sync"
	"testing"

	"github.com/mimecast/dtail/internal"
	"github.com/mimecast/dtail/internal/io/line"
	maprserver "github.com/mimecast/dtail/internal/mapr/server"
	userserver "github.com/mimecast/dtail/internal/user/server"
)

func TestAggregatePointerRaceWithShutdown(t *testing.T) {
	resetServerLogger(t)

	const iterations = 200

	for i := 0; i < iterations; i++ {
		h := &baseHandler{
			done:           internal.NewDone(),
			lines:          make(chan *line.Line, 4),
			serverMessages: make(chan string, 8),
			maprMessages:   make(chan string, 4),
			user:           &userserver.User{Name: "race-test-output-user"},
		}

		ta, err := maprserver.NewAggregate("select count($0) from .", "")
		if err != nil {
			t.Skipf("could not create output aggregate: %v", err)
		}

		var wg sync.WaitGroup

		// Writer goroutine – simulates handleMapCommand setting the output aggregate.
		wg.Add(1)
		go func() {
			defer wg.Done()
			h.setAggregate(ta)
		}()

		// Reader goroutine – simulates Aggregate() / Shutdown reading.
		wg.Add(1)
		go func() {
			defer wg.Done()
			_ = h.getAggregate()
		}()

		wg.Wait()
		// Cleanup: abort the output aggregate so its internal goroutines exit.
		if got := h.getAggregate(); got != nil {
			got.Abort()
		}
	}
}