From 849951be1d1a7ee9f9302006ccb187bf5b4e36f3 Mon Sep 17 00:00:00 2001 From: Paul Buetow Date: Wed, 22 Jul 2026 23:51:18 +0300 Subject: =?UTF-8?q?feat:=20DTail=20fork=20=E2=80=94=20server/client=20feat?= =?UTF-8?q?ure=20development?= MIME-Version: 1.0 Content-Type: text/plain; charset=UTF-8 Content-Transfer-Encoding: 8bit Squashed development of the snonux/dtail fork's product code (internal/, cmd/) since diverging from mimecast/dtail. Major areas: - Read/output path: the former "turbo" channel-less path is now the single, default server-side read/output path for cat/grep/tail and MapReduce; the old channel-based path and its config/env toggles were removed. - MapReduce: single aggregate implementation (server + serverless) fed directly by a processor pipeline, with input-exhausted finalization via the shutdown coordinator; high-concurrency and data-race fixes. - Journal source reads (journal:unit.service) via journalctl, Linux-gated behind a journal-v1 capability. - Auth-key fast reconnect: in-memory per-user public-key cache with TTL/max-keys, registered over an authenticated session (AUTHKEY), checked before authorized_keys. - Interactive query reload (--interactive-query) with SESSION START/UPDATE generation boundaries and capability negotiation. - Client-side deadlines: --timeout / --shutdownAfter as context deadlines; follow shutdown handling. - Client logging: diagnostics-only daily log by default, opt-in payload tee via --log-payload. - Numerous correctness fixes (buffer-pool double-recycle races, EOF-sentinel leaks, glob-expansion cap, TOCTOU in CSV parsing) with accompanying unit tests. Co-Authored-By: Claude Opus 4.8 --- internal/profiling/README.md | 320 ++++++++++++++++++++++++++++++++++++ internal/profiling/flags.go | 38 +++++ internal/profiling/profiler.go | 227 +++++++++++++++++++++++++ internal/profiling/profiler_test.go | 269 ++++++++++++++++++++++++++++++ 4 files changed, 854 insertions(+) create mode 100644 internal/profiling/README.md create mode 100644 internal/profiling/flags.go create mode 100644 internal/profiling/profiler.go create mode 100644 internal/profiling/profiler_test.go (limited to 'internal/profiling') diff --git a/internal/profiling/README.md b/internal/profiling/README.md new file mode 100644 index 0000000..31bdf8a --- /dev/null +++ b/internal/profiling/README.md @@ -0,0 +1,320 @@ +# DTail Internal Profiling Package + +This package provides built-in profiling capabilities for DTail commands (dcat, dgrep, dmap), enabling performance analysis and optimization. + +## Overview + +The profiling package integrates Go's pprof profiling tools directly into DTail commands, allowing users to: +- Profile CPU usage to identify performance bottlenecks +- Track memory allocations to optimize memory usage +- Generate heap profiles to detect memory leaks +- Analyze performance without external profiling tools + +## Architecture + +### Core Components + +1. **Profiler Interface** (`profiler.go`) + - Manages profiling lifecycle (start, stop, write) + - Handles multiple profile types concurrently + - Provides thread-safe operations + +2. **Profile Manager** (`manager.go`) + - Coordinates CPU, memory, and allocation profiling + - Generates timestamped profile filenames + - Manages profile output directories + +3. **Metrics Collection** (`metrics.go`) + - Captures runtime metrics (goroutines, memory stats, GC info) + - Provides before/after snapshots + - Formats metrics for logging + +## Usage + +### Command-Line Flags + +All DTail commands support these profiling flags: + +```bash +# Enable both CPU and memory profiling +dcat -profile -profiledir profiles file.log + +# Enable only CPU profiling +dgrep -cpuprofile -profiledir profiles -regex "error" file.log + +# Enable only memory profiling +dmap -memprofile -profiledir profiles -query "select count(*)" -files data.csv +``` + +### Programmatic Usage + +```go +import "github.com/mimecast/dtail/internal/profiling" + +func main() { + // Create profiler + profiler := profiling.NewProfiler("dcat", "profiles") + + // Start profiling + if err := profiler.Start(); err != nil { + log.Fatal(err) + } + defer profiler.Stop() + + // Your application code here + processFiles() + + // Profiler automatically writes profiles on Stop() +} +``` + +### Advanced Features + +#### Runtime Metrics + +The profiler captures detailed runtime metrics: + +```go +metrics := profiling.GetRuntimeMetrics() +// Includes: NumGoroutines, Alloc, TotalAlloc, Sys, NumGC, GCPauseTotal +``` + +#### Custom Profile Points + +Add profiling snapshots at specific points: + +```go +profiler.LogMetrics("before_processing") +// ... heavy processing ... +profiler.LogMetrics("after_processing") +``` + +## Profile Output + +### File Naming Convention + +Profiles are saved with descriptive filenames: +``` +__.prof + +Examples: +- dcat_cpu_20240626_143022.prof +- dgrep_mem_20240626_143025.prof +- dmap_alloc_20240626_143028.prof +``` + +### Profile Types + +1. **CPU Profile** (`*_cpu_*.prof`) + - Samples CPU usage during execution + - Shows time spent in each function + - Helps identify computational bottlenecks + +2. **Memory Profile** (`*_mem_*.prof`) + - Snapshot of heap allocations at program end + - Shows current memory usage by function + - Useful for finding memory leaks + +3. **Allocation Profile** (`*_alloc_*.prof`) + - Tracks all allocations during execution + - More detailed than memory profile + - Helps reduce allocation pressure + +## Analysis + +### Using pprof + +Analyze profiles with Go's pprof tool: + +```bash +# Interactive mode +go tool pprof profiles/dcat_cpu_20240626_143022.prof + +# Web interface (requires graphviz) +go tool pprof -http=:8080 profiles/dcat_cpu_20240626_143022.prof + +# Generate flame graph +go tool pprof -svg profiles/dgrep_cpu_*.prof > profile.svg + +# Top functions by CPU usage +go tool pprof -top profiles/dmap_cpu_*.prof +``` + +### Using dtail-tools + +DTail provides a convenient tool for profile analysis: + +```bash +# List all profiles +dtail-tools profile -mode list + +# Analyze specific profile +dtail-tools profile -mode analyze profiles/dcat_cpu_*.prof + +# Open web interface +dtail-tools profile -mode analyze profiles/dcat_cpu_*.prof -web +``` + +## Implementation Details + +### Thread Safety + +All profiling operations are thread-safe using sync.Mutex to ensure: +- Safe concurrent access to profiler state +- Atomic start/stop operations +- Protected file writes + +### Error Handling + +The profiler handles errors gracefully: +- Non-blocking: profiling errors don't stop the main program +- Logged warnings for profile write failures +- Automatic cleanup on errors + +### Performance Impact + +Profiling overhead is minimal: +- CPU profiling: ~2-5% overhead +- Memory profiling: Negligible (snapshot at end) +- Allocation profiling: ~5-10% overhead + +## Best Practices + +### When to Profile + +1. **Development Phase** + - Profile new features before merging + - Establish performance baselines + - Identify optimization opportunities + +2. **Performance Issues** + - User reports of slow operations + - High resource usage in production + - Unexplained performance degradation + +3. **Regular Maintenance** + - Periodic profiling to catch regressions + - Before/after optimization work + - Release performance validation + +### Profiling Strategies + +1. **Start Simple** + ```bash + # Quick CPU profile + dcat -cpuprofile -profiledir profiles large_file.log + ``` + +2. **Memory Issues** + ```bash + # Check for memory leaks + dgrep -memprofile -profiledir profiles -regex "pattern" huge_file.log + ``` + +3. **Full Analysis** + ```bash + # Complete profiling suite + dmap -profile -profiledir profiles -query "complex query" -files data.csv + ``` + +### Common Patterns to Look For + +#### CPU Profiles +- Functions with high Flat% (direct CPU usage) +- High cumulative% indicating expensive call chains +- Unexpected hotspots in utility functions + +#### Memory Profiles +- Large allocations in unexpected places +- Many small allocations (consider pooling) +- Growing heap over time (potential leaks) + +## Integration with CI/CD + +### Automated Profiling + +Add profiling to your CI pipeline: + +```yaml +# GitHub Actions example +- name: Run Profiled Tests + run: | + # Run with profiling enabled + go test -bench=. -cpuprofile=cpu.prof -memprofile=mem.prof ./... + + # Analyze results + go tool pprof -top cpu.prof > cpu_analysis.txt + + # Check for regressions + ./scripts/check_performance.sh cpu_analysis.txt +``` + +### Performance Gates + +Set thresholds for key metrics: +```bash +# Example: Fail if function takes >20% CPU +if go tool pprof -top cpu.prof | grep "myFunction" | awk '{print $2}' | grep -E "[2-9][0-9]\.|100\."; then + echo "Performance regression detected!" + exit 1 +fi +``` + +## Troubleshooting + +### No Profiles Generated + +1. Check write permissions for profile directory +2. Ensure command completes successfully +3. Verify profiling flags are correct +4. Check disk space availability + +### Empty Profiles + +1. Increase workload size +2. Ensure program runs long enough +3. Check for early exits or errors +4. Verify profile type matches workload + +### Analysis Errors + +1. Ensure profile format is valid +2. Check Go version compatibility +3. Install graphviz for visualizations +4. Use matching pprof version + +## Contributing + +When modifying the profiling package: + +1. **Maintain Compatibility** + - Don't break existing command-line flags + - Preserve profile output format + - Keep thread-safety guarantees + +2. **Add Tests** + - Unit tests for new functionality + - Integration tests with actual profiles + - Benchmark impact of changes + +3. **Document Changes** + - Update this README + - Add code comments + - Include examples + +## Future Enhancements + +Planned improvements: +- Block profiling support +- Mutex contention profiling +- Custom profile aggregation +- Real-time profiling dashboard +- Profile comparison tools +- Automatic anomaly detection + +## References + +- [Go Profiling Guide](https://go.dev/blog/pprof) +- [pprof Documentation](https://github.com/google/pprof/blob/master/doc/README.md) +- [Runtime Package](https://pkg.go.dev/runtime) +- [Testing Package Profiling](https://pkg.go.dev/testing#hdr-Benchmarks) \ No newline at end of file diff --git a/internal/profiling/flags.go b/internal/profiling/flags.go new file mode 100644 index 0000000..59a6d78 --- /dev/null +++ b/internal/profiling/flags.go @@ -0,0 +1,38 @@ +package profiling + +import "flag" + +// Flags holds command-line flags for profiling +type Flags struct { + // Enable CPU profiling + CPUProfile bool + // Enable memory profiling + MemProfile bool + // Enable all profiling (CPU + memory) + Profile bool + // Directory to store profiles + ProfileDir string +} + +// AddFlags adds profiling flags to the flag set +func AddFlags(f *Flags) { + flag.BoolVar(&f.CPUProfile, "cpuprofile", false, "Enable CPU profiling") + flag.BoolVar(&f.MemProfile, "memprofile", false, "Enable memory profiling") + flag.BoolVar(&f.Profile, "profile", false, "Enable all profiling (CPU + memory)") + flag.StringVar(&f.ProfileDir, "profiledir", "profiles", "Directory to store profiles") +} + +// ToConfig converts flags to profiler config +func (f *Flags) ToConfig(commandName string) Config { + return Config{ + CPUProfile: f.CPUProfile || f.Profile, + MemProfile: f.MemProfile || f.Profile, + ProfileDir: f.ProfileDir, + CommandName: commandName, + } +} + +// Enabled returns true if any profiling is enabled +func (f *Flags) Enabled() bool { + return f.CPUProfile || f.MemProfile || f.Profile +} \ No newline at end of file diff --git a/internal/profiling/profiler.go b/internal/profiling/profiler.go new file mode 100644 index 0000000..8af2567 --- /dev/null +++ b/internal/profiling/profiler.go @@ -0,0 +1,227 @@ +package profiling + +import ( + "fmt" + "os" + "path/filepath" + "runtime" + "runtime/pprof" + "time" + + "log" +) + +// Profiler manages CPU and memory profiling for dtail commands +type Profiler struct { + cpuProfile *os.File + memProfile string + profileDir string + commandName string + enabled bool +} + +// Config holds profiling configuration +type Config struct { + // Enable CPU profiling + CPUProfile bool + // Enable memory profiling + MemProfile bool + // Directory to store profiles + ProfileDir string + // Command name for profile naming + CommandName string +} + +// NewProfiler creates a new profiler instance +func NewProfiler(cfg Config) *Profiler { + if !cfg.CPUProfile && !cfg.MemProfile { + return &Profiler{enabled: false} + } + + p := &Profiler{ + profileDir: cfg.ProfileDir, + commandName: cfg.CommandName, + enabled: true, + } + + // Create profile directory if it doesn't exist + if p.profileDir == "" { + p.profileDir = "profiles" + } + if err := os.MkdirAll(p.profileDir, 0755); err != nil { + log.Printf("Failed to create profile directory: %v", err) + p.enabled = false + return p + } + + // Start CPU profiling if enabled + if cfg.CPUProfile { + p.startCPUProfile() + } + + // Set memory profile path if enabled + if cfg.MemProfile { + timestamp := time.Now().Format("20060102_150405") + p.memProfile = filepath.Join(p.profileDir, fmt.Sprintf("%s_mem_%s.prof", p.commandName, timestamp)) + } + + return p +} + +// startCPUProfile starts CPU profiling +func (p *Profiler) startCPUProfile() { + timestamp := time.Now().Format("20060102_150405") + cpuProfilePath := filepath.Join(p.profileDir, fmt.Sprintf("%s_cpu_%s.prof", p.commandName, timestamp)) + + f, err := os.Create(cpuProfilePath) + if err != nil { + log.Printf("Failed to create CPU profile file: %v", err) + return + } + + if err := pprof.StartCPUProfile(f); err != nil { + log.Printf("Failed to start CPU profile: %v", err) + f.Close() + return + } + + p.cpuProfile = f + log.Printf("Started CPU profiling: %s", cpuProfilePath) +} + +// Stop stops all profiling and writes profiles to disk +func (p *Profiler) Stop() { + if !p.enabled { + return + } + + // Stop CPU profiling + if p.cpuProfile != nil { + pprof.StopCPUProfile() + p.cpuProfile.Close() + log.Printf("Stopped CPU profiling") + } + + // Write memory profile + if p.memProfile != "" { + p.writeMemProfile() + } +} + +// writeMemProfile writes memory allocation profile to disk +func (p *Profiler) writeMemProfile() { + f, err := os.Create(p.memProfile) + if err != nil { + log.Printf("Failed to create memory profile file: %v", err) + return + } + defer f.Close() + + // Force GC before capturing memory profile for more accurate results + runtime.GC() + + if err := pprof.WriteHeapProfile(f); err != nil { + log.Printf("Failed to write memory profile: %v", err) + return + } + + log.Printf("Wrote memory profile: %s", p.memProfile) + + // Also write allocation profile for detailed allocation tracking + allocProfilePath := filepath.Join(p.profileDir, + fmt.Sprintf("%s_alloc_%s.prof", p.commandName, time.Now().Format("20060102_150405"))) + + allocFile, err := os.Create(allocProfilePath) + if err != nil { + log.Printf("Failed to create allocation profile file: %v", err) + return + } + defer allocFile.Close() + + // Set allocation profiling rate to capture more samples + runtime.MemProfileRate = 1 + + if err := pprof.Lookup("allocs").WriteTo(allocFile, 0); err != nil { + log.Printf("Failed to write allocation profile: %v", err) + return + } + + log.Printf("Wrote allocation profile: %s", allocProfilePath) +} + +// Snapshot takes a memory snapshot at any point during execution +func (p *Profiler) Snapshot(label string) { + if !p.enabled || p.memProfile == "" { + return + } + + timestamp := time.Now().Format("20060102_150405") + snapshotPath := filepath.Join(p.profileDir, + fmt.Sprintf("%s_snapshot_%s_%s.prof", p.commandName, label, timestamp)) + + f, err := os.Create(snapshotPath) + if err != nil { + log.Printf("Failed to create snapshot file: %v", err) + return + } + defer f.Close() + + runtime.GC() + if err := pprof.WriteHeapProfile(f); err != nil { + log.Printf("Failed to write snapshot: %v", err) + return + } + + log.Printf("Wrote memory snapshot: %s (label: %s)", snapshotPath, label) +} + +// ProfileMetrics captures and returns current runtime metrics +type ProfileMetrics struct { + // Memory statistics + Alloc uint64 // Bytes allocated and still in use + TotalAlloc uint64 // Bytes allocated (even if freed) + Sys uint64 // Bytes obtained from system + NumGC uint32 // Number of completed GC cycles + LastGC time.Time // Time of last GC + PauseTotalNs uint64 // Total GC pause time in nanoseconds + + // Goroutine count + NumGoroutine int + + // CPU count + NumCPU int +} + +// GetMetrics returns current runtime metrics +func GetMetrics() ProfileMetrics { + var m runtime.MemStats + runtime.ReadMemStats(&m) + + return ProfileMetrics{ + Alloc: m.Alloc, + TotalAlloc: m.TotalAlloc, + Sys: m.Sys, + NumGC: m.NumGC, + LastGC: time.Unix(0, int64(m.LastGC)), + PauseTotalNs: m.PauseTotalNs, + NumGoroutine: runtime.NumGoroutine(), + NumCPU: runtime.NumCPU(), + } +} + +// LogMetrics logs current runtime metrics +func (p *Profiler) LogMetrics(label string) { + if !p.enabled { + return + } + + metrics := GetMetrics() + log.Printf("Profile metrics [%s]: alloc=%.2fMB total_alloc=%.2fMB sys=%.2fMB num_gc=%d gc_pause=%.2fms goroutines=%d", + label, + float64(metrics.Alloc)/1024/1024, + float64(metrics.TotalAlloc)/1024/1024, + float64(metrics.Sys)/1024/1024, + metrics.NumGC, + float64(metrics.PauseTotalNs)/1e6, + metrics.NumGoroutine) +} \ No newline at end of file diff --git a/internal/profiling/profiler_test.go b/internal/profiling/profiler_test.go new file mode 100644 index 0000000..9376611 --- /dev/null +++ b/internal/profiling/profiler_test.go @@ -0,0 +1,269 @@ +package profiling + +import ( + "os" + "path/filepath" + "strings" + "testing" + "time" +) + +func TestProfiler(t *testing.T) { + // Create temporary profile directory + tmpDir := t.TempDir() + + t.Run("DisabledProfiler", func(t *testing.T) { + cfg := Config{ + CPUProfile: false, + MemProfile: false, + ProfileDir: tmpDir, + CommandName: "test", + } + + p := NewProfiler(cfg) + if p.enabled { + t.Error("Profiler should be disabled when no profiling is requested") + } + + // Should not panic + p.Stop() + p.Snapshot("test") + p.LogMetrics("test") + }) + + t.Run("CPUProfileOnly", func(t *testing.T) { + cfg := Config{ + CPUProfile: true, + MemProfile: false, + ProfileDir: tmpDir, + CommandName: "testcpu", + } + + p := NewProfiler(cfg) + if !p.enabled { + t.Error("Profiler should be enabled") + } + + // Do some work to generate CPU samples + doWork(100) + + p.Stop() + + // Check if CPU profile was created + profiles, err := filepath.Glob(filepath.Join(tmpDir, "testcpu_cpu_*.prof")) + if err != nil { + t.Fatalf("Failed to list profiles: %v", err) + } + if len(profiles) == 0 { + t.Error("No CPU profile generated") + } + + // Verify profile exists and has content + for _, profile := range profiles { + info, err := os.Stat(profile) + if err != nil { + t.Errorf("Failed to stat profile %s: %v", profile, err) + } + if info.Size() == 0 { + t.Errorf("Profile %s is empty", profile) + } + } + }) + + t.Run("MemProfileOnly", func(t *testing.T) { + cfg := Config{ + CPUProfile: false, + MemProfile: true, + ProfileDir: tmpDir, + CommandName: "testmem", + } + + p := NewProfiler(cfg) + if !p.enabled { + t.Error("Profiler should be enabled") + } + + // Allocate some memory + allocateMemory() + + p.Stop() + + // Check if memory profiles were created + memProfiles, err := filepath.Glob(filepath.Join(tmpDir, "testmem_mem_*.prof")) + if err != nil { + t.Fatalf("Failed to list memory profiles: %v", err) + } + if len(memProfiles) == 0 { + t.Error("No memory profile generated") + } + + allocProfiles, err := filepath.Glob(filepath.Join(tmpDir, "testmem_alloc_*.prof")) + if err != nil { + t.Fatalf("Failed to list allocation profiles: %v", err) + } + if len(allocProfiles) == 0 { + t.Error("No allocation profile generated") + } + }) + + t.Run("BothProfiles", func(t *testing.T) { + cfg := Config{ + CPUProfile: true, + MemProfile: true, + ProfileDir: tmpDir, + CommandName: "testboth", + } + + p := NewProfiler(cfg) + if !p.enabled { + t.Error("Profiler should be enabled") + } + + // Do work and allocate memory + doWork(100) + allocateMemory() + + p.Stop() + + // Check both profile types + cpuProfiles, _ := filepath.Glob(filepath.Join(tmpDir, "testboth_cpu_*.prof")) + memProfiles, _ := filepath.Glob(filepath.Join(tmpDir, "testboth_mem_*.prof")) + allocProfiles, _ := filepath.Glob(filepath.Join(tmpDir, "testboth_alloc_*.prof")) + + if len(cpuProfiles) == 0 { + t.Error("No CPU profile generated") + } + if len(memProfiles) == 0 { + t.Error("No memory profile generated") + } + if len(allocProfiles) == 0 { + t.Error("No allocation profile generated") + } + }) + + t.Run("Snapshot", func(t *testing.T) { + cfg := Config{ + CPUProfile: false, + MemProfile: true, + ProfileDir: tmpDir, + CommandName: "testsnap", + } + + p := NewProfiler(cfg) + + // Take snapshots + p.Snapshot("before") + allocateMemory() + p.Snapshot("after") + + p.Stop() + + // Check snapshots + snapshots, err := filepath.Glob(filepath.Join(tmpDir, "testsnap_snapshot_*.prof")) + if err != nil { + t.Fatalf("Failed to list snapshots: %v", err) + } + + foundBefore := false + foundAfter := false + for _, snapshot := range snapshots { + if strings.Contains(snapshot, "_before_") { + foundBefore = true + } + if strings.Contains(snapshot, "_after_") { + foundAfter = true + } + } + + if !foundBefore { + t.Error("Before snapshot not found") + } + if !foundAfter { + t.Error("After snapshot not found") + } + }) +} + +func TestGetMetrics(t *testing.T) { + metrics := GetMetrics() + + // Basic sanity checks + if metrics.NumCPU <= 0 { + t.Error("NumCPU should be positive") + } + if metrics.NumGoroutine <= 0 { + t.Error("NumGoroutine should be positive") + } + if metrics.Alloc == 0 { + t.Error("Alloc should not be zero") + } +} + +func TestFlags(t *testing.T) { + f := Flags{} + + // Test default state + if f.Enabled() { + t.Error("Flags should not be enabled by default") + } + + // Test individual flags + f.CPUProfile = true + if !f.Enabled() { + t.Error("Should be enabled when CPUProfile is true") + } + + f.CPUProfile = false + f.MemProfile = true + if !f.Enabled() { + t.Error("Should be enabled when MemProfile is true") + } + + f.MemProfile = false + f.Profile = true + if !f.Enabled() { + t.Error("Should be enabled when Profile is true") + } + + // Test ToConfig + cfg := f.ToConfig("testcmd") + if cfg.CommandName != "testcmd" { + t.Error("CommandName not set correctly") + } + if !cfg.CPUProfile || !cfg.MemProfile { + t.Error("Profile flag should enable both CPU and memory profiling") + } +} + +// Helper functions for testing + +func doWork(iterations int) { + // CPU-intensive work + result := 0 + for i := 0; i < iterations*1000; i++ { + for j := 0; j < 100; j++ { + result += i * j + } + } + _ = result +} + +func allocateMemory() [][]byte { + // Allocate some memory + const numAllocs = 100 + const allocSize = 1024 * 1024 // 1MB + + allocations := make([][]byte, numAllocs) + for i := 0; i < numAllocs; i++ { + allocations[i] = make([]byte, allocSize) + // Touch the memory to ensure it's allocated + for j := 0; j < allocSize; j += 4096 { + allocations[i][j] = byte(i) + } + } + + // Sleep briefly to allow profiler to capture state + time.Sleep(10 * time.Millisecond) + + return allocations +} \ No newline at end of file -- cgit v1.2.3