Skip to content

Commit 42d74ee

Browse files
authored
fix(ci-visibility): harden review edge cases (#4735)
### What does this PR do? Fixes three CI Visibility review findings around logs, settings bootstrap, and shallow Git recovery. This PR makes the CI Visibility logs writer safe when log writes, flushes, payload rotation, and shutdown overlap. The package-level logs state is now synchronized, writer payload rotation is guarded by a writer mutex, uploads are reserved before goroutines start, and log writes accepted before shutdown are flushed exactly once. It also updates the stale msgpack comments for the JSON logs payload path. It also hardens settings initialization so a nil settings response no longer panics. The bootstrap now logs nil settings responses separately from request errors, leaves settings at their zero value, and preserves the existing upload wait behavior. A test seam for client creation lets the nil-response paths be covered directly. Finally, `UnshallowGitRepository` now treats a successful but quiet `git fetch` as success. The HEAD and upstream fetch fallbacks now run only after real command errors, avoiding the previous nil-error panic and unnecessary fallback behavior. ### Motivation Follow-up review found concurrency and nil-handling gaps that could make CI Visibility unstable under parallel log writes, unusual client responses, or quiet Git commands. The highest-risk issue was in logs: the payload object had its own locking, but the writer swapped the payload pointer without a writer-level lock. Concurrent writes, flushes, and stop could race with rotation and lose accepted log entries. The settings and Git fixes address defensive edge cases: nil settings responses should disable features safely, and Git commands that exit successfully do not need stdout to be considered successful. ### Testing - [x] `go test ./internal/civisibility/integrations/logs` - [x] `go test -race ./internal/civisibility/integrations/logs` - [x] `go test ./internal/civisibility/integrations -run 'TestEnsureSettingsInitialization|TestUploadRepositoryChanges|TestSearchCommitsResponse'` - [x] `go test ./internal/civisibility/utils -run 'TestUnshallowGitRepository'` - [x] `go test ./internal/civisibility/...` - [x] `go test -race ./internal/civisibility/...` - [x] `go vet ./internal/civisibility/...` Co-authored-by: tony.redondo <[email protected]>
1 parent 32680c8 commit 42d74ee

13 files changed

Lines changed: 695 additions & 48 deletions

File tree

internal/civisibility/integrations/civisibility_features.go

Lines changed: 22 additions & 4 deletions
Original file line numberDiff line numberDiff line change
@@ -48,9 +48,15 @@ var (
4848
// additionalFeaturesInitializationOnce ensures we do the additional features initialization just once
4949
additionalFeaturesInitializationOnce sync.Once
5050

51+
// additionalFeaturesInitializationMu serializes additional feature initialization with test-only state resets.
52+
additionalFeaturesInitializationMu sync.Mutex
53+
5154
// ciVisibilityRapidClient contains the http rapid client to do CI Visibility queries and upload to the rapid backend
5255
ciVisibilityClient net.Client
5356

57+
// newCIVisibilityClientWithServiceNameFunc creates the CI Visibility client used during settings bootstrap.
58+
newCIVisibilityClientWithServiceNameFunc = net.NewClientWithServiceName
59+
5460
// ciVisibilitySettings contains the CI Visibility settings for this session
5561
ciVisibilitySettings net.SettingsResponseData
5662

@@ -89,7 +95,7 @@ func ensureSettingsInitialization(serviceName string) {
8995
defer log.Debug("civisibility: settings initialization complete")
9096

9197
// Create the CI Visibility client
92-
ciVisibilityClient = net.NewClientWithServiceName(serviceName)
98+
ciVisibilityClient = newCIVisibilityClientWithServiceNameFunc(serviceName)
9399
if ciVisibilityClient == nil {
94100
log.Error("civisibility: error getting the ci visibility http client")
95101
return
@@ -138,7 +144,7 @@ func ensureSettingsInitialization(serviceName string) {
138144
// Get the CI Visibility settings payload for this test session
139145
ciSettings, err := ciVisibilityClient.GetSettings()
140146
if err != nil || ciSettings == nil {
141-
log.Error("civisibility: error getting CI visibility settings: %s", err.Error())
147+
logSettingsFetchError(err)
142148
log.Debug("civisibility: no need to wait for the git upload to finish")
143149
// Enqueue a close action to wait for the upload to finish before finishing the process
144150
PushCiVisibilityCloseAction(waitUploadFactory(time.Minute))
@@ -153,8 +159,8 @@ func ensureSettingsInitialization(serviceName string) {
153159
return
154160
}
155161
ciSettings, err = ciVisibilityClient.GetSettings()
156-
if err != nil {
157-
log.Error("civisibility: error getting CI visibility settings: %s", err.Error())
162+
if err != nil || ciSettings == nil {
163+
logSettingsFetchError(err)
158164
return
159165
}
160166
}
@@ -230,8 +236,20 @@ func ensureSettingsInitialization(serviceName string) {
230236
})
231237
}
232238

239+
// logSettingsFetchError reports a failed or empty CI Visibility settings response.
240+
func logSettingsFetchError(err error) {
241+
if err != nil {
242+
log.Error("civisibility: error getting CI visibility settings: %s", err.Error())
243+
return
244+
}
245+
log.Error("civisibility: error getting CI visibility settings: empty response")
246+
}
247+
233248
// ensureAdditionalFeaturesInitialization loads CI Visibility features that depend on the previously fetched settings.
234249
func ensureAdditionalFeaturesInitialization(_ string) {
250+
additionalFeaturesInitializationMu.Lock()
251+
defer additionalFeaturesInitializationMu.Unlock()
252+
235253
additionalFeaturesInitializationOnce.Do(func() {
236254
log.Debug("civisibility: initializing additional features")
237255
defer log.Debug("civisibility: additional features initialization complete")

internal/civisibility/integrations/civisibility_features_test.go

Lines changed: 126 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -6,10 +6,13 @@
66
package integrations
77

88
import (
9+
"io"
910
"testing"
1011

1112
"github.com/stretchr/testify/assert"
1213
"github.com/stretchr/testify/require"
14+
15+
civisibilitynet "github.com/DataDog/dd-trace-go/v2/internal/civisibility/utils/net"
1316
)
1417

1518
func TestSearchCommitsResponseMissingCommitsPreservesLocalOrder(t *testing.T) {
@@ -99,3 +102,126 @@ func TestUploadRepositoryChangesReusesInitialMissingCommitsWhenUnshallowIsUnavai
99102
assert.Equal(t, []string{"local-1", "local-2"}, uploadedIncludes)
100103
assert.Equal(t, []string{"remote-1"}, uploadedExcludes)
101104
}
105+
106+
func TestEnsureSettingsInitializationNilClientFactoryDoesNotStartUpload(t *testing.T) {
107+
resetCIVisibilityStateForTesting()
108+
t.Cleanup(resetCIVisibilityStateForTesting)
109+
110+
newCIVisibilityClientWithServiceNameFunc = func(_ string) civisibilitynet.Client {
111+
return nil
112+
}
113+
uploadRepositoryChangesFunc = func() (int64, error) {
114+
t.Fatal("repository upload should not start without a CI Visibility client")
115+
return 0, nil
116+
}
117+
118+
require.NotPanics(t, func() {
119+
ensureSettingsInitialization("service")
120+
})
121+
assert.Equal(t, civisibilitynet.SettingsResponseData{}, ciVisibilitySettings)
122+
assert.Len(t, closeActions, 0)
123+
}
124+
125+
func TestEnsureSettingsInitializationHandlesNilInitialSettingsResponse(t *testing.T) {
126+
resetCIVisibilityStateForTesting()
127+
t.Cleanup(resetCIVisibilityStateForTesting)
128+
129+
uploadStarted := make(chan struct{})
130+
uploadRelease := make(chan struct{})
131+
newCIVisibilityClientWithServiceNameFunc = func(_ string) civisibilitynet.Client {
132+
return &mockCIVisibilityClient{
133+
getSettings: func() (*civisibilitynet.SettingsResponseData, error) {
134+
return nil, nil
135+
},
136+
}
137+
}
138+
uploadRepositoryChangesFunc = func() (int64, error) {
139+
close(uploadStarted)
140+
<-uploadRelease
141+
return 0, nil
142+
}
143+
144+
require.NotPanics(t, func() {
145+
ensureSettingsInitialization("service")
146+
})
147+
assert.Equal(t, civisibilitynet.SettingsResponseData{}, ciVisibilitySettings)
148+
assert.Len(t, closeActions, 1)
149+
150+
close(uploadRelease)
151+
closeActions[0]()
152+
<-uploadStarted
153+
}
154+
155+
func TestEnsureSettingsInitializationHandlesNilRetrySettingsResponse(t *testing.T) {
156+
resetCIVisibilityStateForTesting()
157+
t.Cleanup(resetCIVisibilityStateForTesting)
158+
159+
settingsCalls := 0
160+
newCIVisibilityClientWithServiceNameFunc = func(_ string) civisibilitynet.Client {
161+
return &mockCIVisibilityClient{
162+
getSettings: func() (*civisibilitynet.SettingsResponseData, error) {
163+
settingsCalls++
164+
if settingsCalls == 1 {
165+
return &civisibilitynet.SettingsResponseData{RequireGit: true}, nil
166+
}
167+
return nil, nil
168+
},
169+
}
170+
}
171+
uploadRepositoryChangesFunc = func() (int64, error) {
172+
return 0, nil
173+
}
174+
175+
require.NotPanics(t, func() {
176+
ensureSettingsInitialization("service")
177+
})
178+
assert.Equal(t, 2, settingsCalls)
179+
assert.Equal(t, civisibilitynet.SettingsResponseData{}, ciVisibilitySettings)
180+
assert.Len(t, closeActions, 0)
181+
}
182+
183+
// mockCIVisibilityClient implements net.Client for settings bootstrap tests.
184+
type mockCIVisibilityClient struct {
185+
getSettings func() (*civisibilitynet.SettingsResponseData, error)
186+
}
187+
188+
var _ civisibilitynet.Client = (*mockCIVisibilityClient)(nil)
189+
190+
func (m *mockCIVisibilityClient) GetSettings() (*civisibilitynet.SettingsResponseData, error) {
191+
if m.getSettings != nil {
192+
return m.getSettings()
193+
}
194+
return &civisibilitynet.SettingsResponseData{}, nil
195+
}
196+
197+
func (m *mockCIVisibilityClient) GetKnownTests() (*civisibilitynet.KnownTestsResponseData, error) {
198+
return nil, nil
199+
}
200+
201+
func (m *mockCIVisibilityClient) GetCommits(_ []string) ([]string, error) {
202+
return nil, nil
203+
}
204+
205+
func (m *mockCIVisibilityClient) SendPackFiles(_ string, _ []string) (int64, error) {
206+
return 0, nil
207+
}
208+
209+
func (m *mockCIVisibilityClient) SendCoveragePayload(_ io.Reader) error {
210+
return nil
211+
}
212+
213+
func (m *mockCIVisibilityClient) SendCoveragePayloadWithFormat(_ io.Reader, _ string) error {
214+
return nil
215+
}
216+
217+
func (m *mockCIVisibilityClient) GetSkippableTests() (string, map[string]map[string][]civisibilitynet.SkippableResponseDataAttributes, error) {
218+
return "", nil, nil
219+
}
220+
221+
func (m *mockCIVisibilityClient) GetTestManagementTests() (*civisibilitynet.TestManagementTestsResponseDataModules, error) {
222+
return nil, nil
223+
}
224+
225+
func (m *mockCIVisibilityClient) SendLogs(_ io.Reader) error {
226+
return nil
227+
}

internal/civisibility/integrations/gotesting/failnow/failnow_test.go

Lines changed: 8 additions & 0 deletions
Original file line numberDiff line numberDiff line change
@@ -80,6 +80,14 @@ func TestCleanupRunsAfterParallelSubtest(t *testing.T) {
8080
runTestScenario(t, "test-cleanup-after-parallel-subtest", "^TestCleanupRunsAfterParallelSubtestFixture$")
8181
}
8282

83+
// TestParallelSubtestSchedulerSlotIsReleased verifies Datadog-managed retry
84+
// clones release their parent scheduler slot before waiting for parallel
85+
// subtests. With -parallel=1, failing to release that slot deadlocks the child
86+
// in testing.(*testState).waitParallel until the package timeout fires.
87+
func TestParallelSubtestSchedulerSlotIsReleased(t *testing.T) {
88+
runSubprocess(t, "test-cleanup-after-parallel-subtest", "-test.run", "^TestCleanupRunsAfterParallelSubtestFixture$", "-test.parallel=1", "-test.timeout=5s")
89+
}
90+
8391
func TestFlakyRetryGlobalBudget(t *testing.T) {
8492
runTestScenario(t, "test-flaky-retry-global-budget", "^TestFlakyRetryGlobalBudgetFixture$")
8593
}

internal/civisibility/integrations/gotesting/instrumentation.go

Lines changed: 34 additions & 5 deletions
Original file line numberDiff line numberDiff line change
@@ -833,7 +833,7 @@ func executeTestIteration(execOpts *executionOptions) bool {
833833
*cn <- struct{}{}
834834
}()
835835
defer func() {
836-
completeParallelSubtests(localTPrivateFields)
836+
completeParallelSubtests(pLocalT, localTPrivateFields)
837837
}()
838838
defer func() {
839839
duration = time.Since(startTime)
@@ -912,7 +912,7 @@ func executeTestIteration(execOpts *executionOptions) bool {
912912
// cleanup Goexit in a helper goroutine so retry orchestration can treat cleanup
913913
// failures as attempt failures instead of letting them escape the retry loop.
914914
func runTestCleanup(t *testing.T, result *testCleanupResult) {
915-
completeParallelSubtests(getTestPrivateFields(t))
915+
completeParallelSubtests(t, getTestPrivateFields(t))
916916
result.ran = true
917917
done := make(chan struct{})
918918
go func() {
@@ -933,15 +933,20 @@ func runTestCleanup(t *testing.T, result *testCleanupResult) {
933933
}
934934

935935
// completeParallelSubtests releases and waits for parallel subtests owned by a
936-
// Datadog-managed clone. It clears the subtest queue before releasing the
937-
// barrier so later cleanup paths cannot close the same barrier twice.
938-
func completeParallelSubtests(localTPrivateFields *commonPrivateFields) {
936+
// Datadog-managed clone. It mirrors testing.tRunner's scheduler accounting:
937+
// release the parent slot before unblocking children, then reacquire it for
938+
// sequential parents before running cleanup.
939+
func completeParallelSubtests(t *testing.T, localTPrivateFields *commonPrivateFields) {
939940
if localTPrivateFields == nil || localTPrivateFields.sub == nil || len(*localTPrivateFields.sub) == 0 {
940941
return
941942
}
942943

943944
subtests := *localTPrivateFields.sub
944945
*localTPrivateFields.sub = nil
946+
testState := getTestState(t)
947+
if testState != nil {
948+
testingTestStateRelease(testState)
949+
}
945950
if localTPrivateFields.barrier != nil && *localTPrivateFields.barrier != nil {
946951
close(*localTPrivateFields.barrier)
947952
}
@@ -951,6 +956,24 @@ func completeParallelSubtests(localTPrivateFields *commonPrivateFields) {
951956
<-*pvSub.signal
952957
}
953958
}
959+
if testState != nil && !isParallelTest(t, localTPrivateFields) {
960+
testingTestStateWaitParallel(testState)
961+
}
962+
}
963+
964+
// isParallelTest reports whether the active test has entered Go's parallel-test
965+
// path. Datadog-managed retry clones forward Parallel to the original *testing.T,
966+
// so the original must also be checked before deciding whether to reacquire the
967+
// scheduler slot.
968+
func isParallelTest(t *testing.T, localTPrivateFields *commonPrivateFields) bool {
969+
if localTPrivateFields != nil && localTPrivateFields.isParallel != nil && *localTPrivateFields.isParallel {
970+
return true
971+
}
972+
if execMeta := getTestMetadata(t); execMeta != nil && execMeta.originalTest != nil {
973+
originalFields := getTestPrivateFields(execMeta.originalTest)
974+
return originalFields != nil && originalFields.isParallel != nil && *originalFields.isParallel
975+
}
976+
return false
954977
}
955978

956979
// runAndApplyTestCleanup runs a retry attempt's cleanups before its span is
@@ -1019,3 +1042,9 @@ func (m *noopMutex) TryLock() bool { return true }
10191042

10201043
//go:linkname testingTRunCleanup testing.(*common).runCleanup
10211044
func testingTRunCleanup(c *testing.T, ph int) (panicVal any)
1045+
1046+
//go:linkname testingTestStateWaitParallel testing.(*testState).waitParallel
1047+
func testingTestStateWaitParallel(s *testingTestState)
1048+
1049+
//go:linkname testingTestStateRelease testing.(*testState).release
1050+
func testingTestStateRelease(s *testingTestState)

internal/civisibility/integrations/gotesting/reflections.go

Lines changed: 51 additions & 20 deletions
Original file line numberDiff line numberDiff line change
@@ -78,16 +78,17 @@ func copyFieldUsingPointersWithConversion[V any](source any, target any, fieldNa
7878

7979
// commonPrivateFields is collection of required private fields from testing.common
8080
type commonPrivateFields struct {
81-
mu *sync.RWMutex
82-
output *[]byte // Output generated by test or benchmark.
83-
level *int
84-
name *string // Name of test or benchmark.
85-
failed *bool // Test or benchmark has failed.
86-
skipped *bool // Test or benchmark has been skipped.
87-
parent *unsafe.Pointer // Parent common
88-
barrier *chan bool // Barrier for parallel tests
89-
signal *chan bool // Signal channel for test completion
90-
sub *[]*testing.T // Queue of subtests to be run in parallel.
81+
mu *sync.RWMutex
82+
output *[]byte // Output generated by test or benchmark.
83+
level *int
84+
name *string // Name of test or benchmark.
85+
failed *bool // Test or benchmark has failed.
86+
skipped *bool // Test or benchmark has been skipped.
87+
isParallel *bool // Whether the test has called testing.T.Parallel.
88+
parent *unsafe.Pointer // Parent common.
89+
barrier *chan bool // Barrier for parallel tests.
90+
signal *chan bool // Signal channel for test completion.
91+
sub *[]*testing.T // Queue of subtests to be run in parallel.
9192
}
9293

9394
// AddLevel increase or decrease the testing.common.level field value, used by
@@ -188,16 +189,17 @@ func getTestPrivateFieldsFast(t *testing.T, layout *testingInternalsLayout) *com
188189
return nil
189190
}
190191
fields := &commonPrivateFields{
191-
mu: fieldPtr[sync.RWMutex](commonBase, layout.common.mu),
192-
output: fieldPtr[[]byte](commonBase, layout.common.output),
193-
level: fieldPtr[int](commonBase, layout.common.level),
194-
name: fieldPtr[string](commonBase, layout.common.name),
195-
failed: fieldPtr[bool](commonBase, layout.common.failed),
196-
skipped: fieldPtr[bool](commonBase, layout.common.skipped),
197-
parent: (*unsafe.Pointer)(fieldRawPtr(commonBase, layout.common.parent.unsafeField)),
198-
barrier: fieldPtr[chan bool](commonBase, layout.common.barrier),
199-
signal: fieldPtr[chan bool](commonBase, layout.common.signal),
200-
sub: fieldPtr[[]*testing.T](commonBase, layout.common.sub),
192+
mu: fieldPtr[sync.RWMutex](commonBase, layout.common.mu),
193+
output: fieldPtr[[]byte](commonBase, layout.common.output),
194+
level: fieldPtr[int](commonBase, layout.common.level),
195+
name: fieldPtr[string](commonBase, layout.common.name),
196+
failed: fieldPtr[bool](commonBase, layout.common.failed),
197+
skipped: fieldPtr[bool](commonBase, layout.common.skipped),
198+
isParallel: fieldPtr[bool](commonBase, layout.common.isParallel),
199+
parent: (*unsafe.Pointer)(fieldRawPtr(commonBase, layout.common.parent.unsafeField)),
200+
barrier: fieldPtr[chan bool](commonBase, layout.common.barrier),
201+
signal: fieldPtr[chan bool](commonBase, layout.common.signal),
202+
sub: fieldPtr[[]*testing.T](commonBase, layout.common.sub),
201203
}
202204
runtime.KeepAlive(t)
203205
return fields
@@ -227,6 +229,9 @@ func getTestPrivateFieldsReflect(t *testing.T) *commonPrivateFields {
227229
if ptr, err := getFieldPointerFrom(t, "skipped"); err == nil && ptr != nil {
228230
testFields.skipped = (*bool)(ptr)
229231
}
232+
if ptr, err := getFieldPointerFrom(t, "isParallel"); err == nil && ptr != nil {
233+
testFields.isParallel = (*bool)(ptr)
234+
}
230235
if ptr, err := getFieldPointerFrom(t, "parent"); err == nil && ptr != nil {
231236
testFields.parent = (*unsafe.Pointer)(ptr)
232237
}
@@ -243,6 +248,32 @@ func getTestPrivateFieldsReflect(t *testing.T) *commonPrivateFields {
243248
return testFields
244249
}
245250

251+
// testingTestState is an opaque handle for testing.testState. The real type is
252+
// private to the standard library; callers must only pass values obtained from
253+
// a *testing.T's private tstate field.
254+
type testingTestState struct{}
255+
256+
// getTestState returns the private testing.testState pointer used by Go's
257+
// parallel-test scheduler. A nil result means this Go runtime layout is not
258+
// supported by the scheduler-slot helper.
259+
func getTestState(t *testing.T) *testingTestState {
260+
layout := getTestingInternalsLayout()
261+
if layout != nil && !layout.disabled && layout.tstate.available {
262+
ptr := fieldRawPtr(unsafe.Pointer(t), layout.tstate.unsafeField)
263+
if ptr == nil {
264+
return nil
265+
}
266+
state := *(**testingTestState)(ptr)
267+
runtime.KeepAlive(t)
268+
return state
269+
}
270+
271+
if ptr, err := getFieldPointerFrom(t, "tstate"); err == nil && ptr != nil {
272+
return *(**testingTestState)(ptr)
273+
}
274+
return nil
275+
}
276+
246277
// getTestParentPrivateFields is a method to retrieve all required parent privates field from
247278
// testing.T.parent, returning a commonPrivateFields instance
248279
func getTestParentPrivateFields(t *testing.T) *commonPrivateFields {

0 commit comments

Comments
 (0)