diff --git a/integration/agentcompat/internal/scenario/stress_fd_diagnostic_collector_test.go b/integration/agentcompat/internal/scenario/stress_fd_diagnostic_collector_test.go new file mode 100644 index 00000000..d4980abf --- /dev/null +++ b/integration/agentcompat/internal/scenario/stress_fd_diagnostic_collector_test.go @@ -0,0 +1,81 @@ +//go:build linux && agentcompat + +package scenario + +import ( + "encoding/json" + "fmt" + "strings" + "testing" + + "github.com/stretchr/testify/require" + + processharness "github.com/nezhahq/nezha/integration/agentcompat/internal/process" +) + +func TestFDDiagnosticCollector_UsesOnlyFinalSampleCountThroughCollectorWiring(t *testing.T) { + // Given + collector := newFDDiagnosticCollector(fdDiagnosticCollectorSpec{Enabled: true, TailResultCapacity: 8, Sample: fdDiagnosticTestSampler}) + pair := stressDiagnosticAgentWindow(t, stressDiagnosticAgentWindowSpec{Ordinal: 1, PID: 101, BaselineCount: 8, EndCount: 8, Target: "stable"}) + pair.end.Window.Samples[0].NonStdioFDCount = 9 + collector.RecordBaseline(pair.baseline) + + // When + collector.RecordEnd(t.Context(), pair.end) + records := collector.WaitRecords() + + // Then + require.Empty(t, records) +} + +func TestFDDiagnosticCollector_OrdersAndLogsStableCompleteRecords(t *testing.T) { + // Given + collector := newFDDiagnosticCollector(fdDiagnosticCollectorSpec{Enabled: true, SamplerPID: 900, TailResultCapacity: 8, Sample: fdDiagnosticTestSampler}) + agentOne := stressDiagnosticAgentWindow(t, stressDiagnosticAgentWindowSpec{Ordinal: 1, PID: 101, BaselineCount: 8, EndCount: 9, Target: "one"}) + agentTwo := stressDiagnosticAgentWindow(t, stressDiagnosticAgentWindowSpec{Ordinal: 2, PID: 102, BaselineCount: 8, EndCount: 9, Target: "two"}) + fdDiagnosticSetFinalObservations(&agentOne, 101) + fdDiagnosticSetFinalObservations(&agentTwo, 102) + collector.RecordBaseline(agentTwo.baseline) + collector.RecordBaseline(agentOne.baseline) + + // When + collector.RecordEnd(t.Context(), agentTwo.end) + collector.RecordEnd(t.Context(), agentOne.end) + records := collector.WaitRecords() + logs := make([]string, 0, len(records)) + logger := &fdDiagnosticMemoryLogger{lines: &logs} + collector.WaitAndLog(logger) + collector.WaitAndLog(logger) + + // Then + require.Len(t, records, 2) + require.Equal(t, []int{1, 2}, []int{records[0].AgentOrdinal, records[1].AgentOrdinal}) + require.Len(t, logs, 2) + for index, line := range logs { + require.True(t, strings.HasPrefix(line, "agentcompat_fd_diagnostic=")) + var record fdDiagnosticRecord + require.NoError(t, json.Unmarshal([]byte(strings.TrimPrefix(line, "agentcompat_fd_diagnostic=")), &record)) + require.Equal(t, index+1, record.AgentOrdinal) + require.NotEmpty(t, record.Baseline.Samples[0].FDObservations) + require.NotEmpty(t, record.End.Samples[0].FDObservations) + require.NotEmpty(t, record.Tail[0].FDObservations) + require.False(t, record.Baseline.Samples[0].ObservedAt.IsZero()) + require.False(t, record.End.Samples[0].ObservedAt.IsZero()) + require.False(t, record.Tail[0].ObservedAt.IsZero()) + require.Equal(t, "observed_through_tail", record.Lifecycle[0].Status) + } +} + +func fdDiagnosticSetFinalObservations(pair *fdDiagnosticWindowPair, pid int) { + lastSample := len(pair.baseline.Window.Samples) - 1 + pair.baseline.Window.Samples[lastSample].FDObservations = []processharness.FDObservation{{Number: 3, Target: fmt.Sprintf("baseline-%d", pid)}} + pair.end.Window.Samples[lastSample].FDObservations = []processharness.FDObservation{{Number: 4, Target: fmt.Sprintf("added-%d", pid)}} +} + +type fdDiagnosticMemoryLogger struct { + lines *[]string +} + +func (logger *fdDiagnosticMemoryLogger) Logf(format string, arguments ...any) { + *logger.lines = append(*logger.lines, fmt.Sprintf(format, arguments...)) +} diff --git a/integration/agentcompat/internal/scenario/stress_fd_diagnostic_tail_test.go b/integration/agentcompat/internal/scenario/stress_fd_diagnostic_tail_test.go new file mode 100644 index 00000000..3bea2e83 --- /dev/null +++ b/integration/agentcompat/internal/scenario/stress_fd_diagnostic_tail_test.go @@ -0,0 +1,160 @@ +//go:build linux && agentcompat + +package scenario + +import ( + "context" + "sync" + "testing" + "time" + + "github.com/stretchr/testify/require" + + processharness "github.com/nezhahq/nezha/integration/agentcompat/internal/process" +) + +func TestFDDiagnosticCollector_StartsOverlappingTailsInOrdinalOrder(t *testing.T) { + // Given + ctx, cancel := context.WithCancel(t.Context()) + defer cancel() + firstStarts := make(chan int, 2) + secondStarts := make(chan int, 2) + completed := make(chan int, 2) + release := map[int]chan struct{}{101: make(chan struct{}), 102: make(chan struct{})} + counts := make(map[int]int) + var countsMu sync.Mutex + collector := newFDDiagnosticCollector(fdDiagnosticCollectorSpec{ + Enabled: true, + SamplerPID: 900, + TailResultCapacity: 8, + TailInterval: 0, + Sample: func(ctx context.Context, pid int) (processharness.Sample, error) { + countsMu.Lock() + counts[pid]++ + ordinal := counts[pid] + countsMu.Unlock() + if ordinal == 1 { + firstStarts <- pid + } + if ordinal == 2 { + secondStarts <- pid + select { + case <-release[pid]: + case <-ctx.Done(): + return processharness.Sample{}, ctx.Err() + } + } + if ordinal == 20 { + completed <- pid + } + return processharness.Sample{PID: pid, FDObservations: []processharness.FDObservation{{Number: 3, Target: "tail"}}, SampledAt: time.Unix(int64(ordinal), 0).UTC()}, nil + }, + }) + pairOne := stressDiagnosticAgentWindow(t, stressDiagnosticAgentWindowSpec{Ordinal: 1, PID: 101, BaselineCount: 8, EndCount: 9, Target: "one"}) + pairTwo := stressDiagnosticAgentWindow(t, stressDiagnosticAgentWindowSpec{Ordinal: 2, PID: 102, BaselineCount: 8, EndCount: 9, Target: "two"}) + collector.RecordBaseline(pairOne.baseline) + collector.RecordBaseline(pairTwo.baseline) + + // When + collector.RecordEnd(ctx, pairOne.end) + require.Equal(t, 101, <-firstStarts) + require.Equal(t, 101, <-secondStarts) + collector.RecordEnd(ctx, pairTwo.end) + require.Equal(t, 102, <-firstStarts) + require.Equal(t, 102, <-secondStarts) + close(release[102]) + require.Equal(t, 102, <-completed) + close(release[101]) + require.Equal(t, 101, <-completed) + records := collector.WaitRecords() + + // Then + require.Len(t, records, 2) + require.Equal(t, 1, records[0].AgentOrdinal) + require.Equal(t, 2, records[1].AgentOrdinal) + for _, record := range records { + require.Len(t, record.Tail, 20) + require.Equal(t, 1, record.Tail[0].Ordinal) + require.Equal(t, 20, record.Tail[19].Ordinal) + require.Equal(t, time.Unix(1, 0).UTC(), record.Tail[0].ObservedAt) + require.Equal(t, time.Unix(20, 0).UTC(), record.Tail[19].ObservedAt) + require.Equal(t, "tail_complete", record.LifecycleStatus) + } +} + +func TestFDDiagnosticCollector_DoesNotSampleWhenDisabled(t *testing.T) { + // Given + called := false + collector := newFDDiagnosticCollector(fdDiagnosticCollectorSpec{ + Enabled: false, + Sample: func(context.Context, int) (processharness.Sample, error) { + called = true + return processharness.Sample{}, nil + }, + }) + pair := stressDiagnosticAgentWindow(t, stressDiagnosticAgentWindowSpec{Ordinal: 1, PID: 101, BaselineCount: 8, EndCount: 9, Target: "disabled"}) + + // When + collector.RecordBaseline(pair.baseline) + collector.RecordEnd(t.Context(), pair.end) + records := collector.WaitRecords() + + // Then + require.False(t, called) + require.Empty(t, records) +} + +func TestFDDiagnosticCollector_NilReceiverIsDisabled(t *testing.T) { + // Given + var collector *fdDiagnosticCollector + pair := stressDiagnosticAgentWindow(t, stressDiagnosticAgentWindowSpec{Ordinal: 1, PID: 101, BaselineCount: 8, EndCount: 9, Target: "nil"}) + + // When / Then + require.NotPanics(t, func() { + require.False(t, collector.Enabled()) + collector.RecordBaseline(pair.baseline) + collector.RecordEnd(t.Context(), pair.end) + require.Empty(t, collector.WaitRecords()) + }) +} + +func TestFDDiagnosticTail_UsesConfiguredTimerIntervalAfterFirstSample(t *testing.T) { + // Given + ctx, cancel := context.WithCancel(t.Context()) + defer cancel() + firstSample := make(chan struct{}, 1) + sampler := func(context.Context, int) (processharness.Sample, error) { + firstSample <- struct{}{} + return processharness.Sample{PID: 101}, nil + } + + // When + result := make(chan fdDiagnosticTailResult, 1) + go func() { + result <- collectFDDiagnosticTail(fdDiagnosticTailSpec{Context: ctx, PID: 101, Interval: time.Hour, Sample: sampler}) + }() + <-firstSample + cancel() + tail := <-result + + // Then + require.Len(t, tail.Samples, 1) + require.ErrorIs(t, tail.Err, context.Canceled) +} + +func TestFDDiagnosticTail_DoesNotInvokeSamplerAfterCancellation(t *testing.T) { + // Given + ctx, cancel := context.WithCancel(t.Context()) + cancel() + called := false + + // When + tail := collectFDDiagnosticTail(fdDiagnosticTailSpec{Context: ctx, PID: 101, Interval: time.Hour, Sample: func(context.Context, int) (processharness.Sample, error) { + called = true + return processharness.Sample{}, nil + }}) + + // Then + require.False(t, called) + require.ErrorIs(t, tail.Err, context.Canceled) +}