This is an automated email from the ASF dual-hosted git repository.

tvalentyn pushed a commit to branch master
in repository https://gitbox.apache.org/repos/asf/beam.git


The following commit(s) were added to refs/heads/master by this push:
     new 8b4c6751963 Support core dump analysis with pystack and gdb. (#39484)
8b4c6751963 is described below

commit 8b4c6751963d14f574bad011ddbaff15ab013a57
Author: tvalentyn <[email protected]>
AuthorDate: Mon Aug 3 05:59:26 2026 -0700

    Support core dump analysis with pystack and gdb. (#39484)
    
    * Refactor process tracking and signal handling in Python container 
bootloader
    
    * Support collecting and processing dumped core files with pystack.
    
    * Terminate postprocessing after profiler was disengaged.
    
    * Also save core analysis in text files.
    
    * Support gdb
    
    * Log gdb commands.
    
    * Change prefix
    
    * Respect the postprocessing interval supplied in options.
    
    * Register a callback to do a final profile postprocessing before crashing 
the container.
    
    * Wait until core file creation finishes.
---
 CHANGES.md                                         |   2 +-
 .../python/apache_beam/options/pipeline_options.py |   4 +
 .../apache_beam/options/pipeline_options_test.py   |   9 +
 .../container/base_image_requirements_manual.txt   |   1 +
 sdks/python/container/boot.go                      |  55 ++--
 .../container/ml/py310/base_image_requirements.txt |   1 +
 .../container/ml/py310/gpu_image_requirements.txt  |   1 +
 .../container/ml/py311/base_image_requirements.txt |   1 +
 .../container/ml/py311/gpu_image_requirements.txt  |   1 +
 .../container/ml/py312/base_image_requirements.txt |   1 +
 .../container/ml/py312/gpu_image_requirements.txt  |   1 +
 .../container/ml/py313/base_image_requirements.txt |   1 +
 sdks/python/container/profiler.go                  | 303 +++++++++++++++++++--
 sdks/python/container/profiler_test.go             | 125 +++++++++
 .../container/py310/base_image_requirements.txt    |   1 +
 .../container/py311/base_image_requirements.txt    |   1 +
 .../container/py312/base_image_requirements.txt    |   1 +
 .../container/py313/base_image_requirements.txt    |   1 +
 .../container/py314/base_image_requirements.txt    |   1 +
 19 files changed, 466 insertions(+), 45 deletions(-)

diff --git a/CHANGES.md b/CHANGES.md
index d853314a0ad..076413d2043 100644
--- a/CHANGES.md
+++ b/CHANGES.md
@@ -76,7 +76,7 @@
 * (Python) Added `equal_to_approx`, an `assert_that` matcher that compares 
numeric pipeline outputs with a configurable tolerance 
([#18028](https://github.com/apache/beam/issues/18028)).
 * (Python) `Timestamp` now supports variable subsecond precision, up to 
nanoseconds. The portable
   `beam:logical_type:timestamp:v1` logical type now maps to Python's 
`Timestamp` ([#39344](https://github.com/apache/beam/issues/39344)).
-* X feature added (Java/Python) 
([#X](https://github.com/apache/beam/issues/X)).
+* (Python) Added support to analyze core dumps created after python worker 
segmentation faults with `pystack` (or `gdb` if installed) using the 
`--profiler_agent=coredump` pipeline option. 
([#39484](https://github.com/apache/beam/issues/39484)).
 
 ## Breaking Changes
 
diff --git a/sdks/python/apache_beam/options/pipeline_options.py 
b/sdks/python/apache_beam/options/pipeline_options.py
index 239d577cfa0..ee7e14f3de2 100644
--- a/sdks/python/apache_beam/options/pipeline_options.py
+++ b/sdks/python/apache_beam/options/pipeline_options.py
@@ -1743,6 +1743,10 @@ class ProfilingOptions(PipelineOptions):
           _LOGGER.info(
               'Setting --profile_location to %s since profiling is enabled.',
               self.profile_location)
+
+      if self.profiler_agent == 'coredump':
+        debug_options = self.view_as(DebugOptions)
+        debug_options.add_experiment('core_pattern=/tmp/beam_coredump.%e.%p')
     return errors
 
 
diff --git a/sdks/python/apache_beam/options/pipeline_options_test.py 
b/sdks/python/apache_beam/options/pipeline_options_test.py
index b321314b4b7..fbdaf25f0e8 100644
--- a/sdks/python/apache_beam/options/pipeline_options_test.py
+++ b/sdks/python/apache_beam/options/pipeline_options_test.py
@@ -702,6 +702,15 @@ class PipelineOptionsTest(unittest.TestCase):
     self.assertTrue(
         any('--profiler_agent is mutually exclusive' in err for err in errors))
 
+  def test_profiling_agent_coredump_adds_core_pattern(self):
+    options = PipelineOptions(['--profiler_agent=coredump'])
+    validator = PipelineOptionsValidator(options, None)
+    self.assertEqual(validator.validate(), [])
+    debug_options = options.view_as(DebugOptions)
+    self.assertEqual(
+        debug_options.lookup_experiment('core_pattern'),
+        '/tmp/beam_coredump.%e.%p')
+
   def test_profile_location_defaulting_and_opt_out(self):
     options = PipelineOptions(
         ['--profiler_agent=memray', '--temp_location=gs://bucket/temp'])
diff --git a/sdks/python/container/base_image_requirements_manual.txt 
b/sdks/python/container/base_image_requirements_manual.txt
index a78d993461c..f771b66ee6d 100644
--- a/sdks/python/container/base_image_requirements_manual.txt
+++ b/sdks/python/container/base_image_requirements_manual.txt
@@ -42,6 +42,7 @@ guppy3
 memray==1.19.3
 mmh3  # Optimizes execution of some Beam codepaths. TODO: Make it Beam's 
dependency.
 nltk  # Commonly used for natural language processing.
+pystack
 google-crc32c
 scipy
 scikit-learn
diff --git a/sdks/python/container/boot.go b/sdks/python/container/boot.go
index 364a614b4e8..5a8d6da46ab 100644
--- a/sdks/python/container/boot.go
+++ b/sdks/python/container/boot.go
@@ -60,6 +60,9 @@ var (
        provisionEndpoint = flag.String("provision_endpoint", "", "Provision 
endpoint (required).")
        controlEndpoint   = flag.String("control_endpoint", "", "Control 
endpoint (required).")
        semiPersistDir    = flag.String("semi_persist_dir", "/tmp", "Local 
semi-persistent directory (optional).")
+
+       workerMu     sync.Mutex
+       shuttingDown bool
 )
 
 const (
@@ -307,19 +310,12 @@ func launchSDKProcess() error {
 
        workerIds := append([]string{*id}, info.GetSiblingWorkerIds()...)
 
-       // Keep track of child PIDs for clean shutdown without zombies
-       childPids := struct {
-               v        []int
-               canceled bool
-               mu       sync.Mutex
-       }{v: make([]int, 0, len(workerIds))}
-
        // Forward trapped signals to child process groups in order to 
terminate them gracefully and avoid zombies
        go func() {
                logger.Printf(ctx, "Received signal: %v", <-signalChannel)
-               childPids.mu.Lock()
-               childPids.canceled = true
-               for _, pid := range childPids.v {
+               workerMu.Lock()
+               shuttingDown = true
+               for _, pid := range activePids {
                        go func(pid int) {
                                // This goroutine will be canceled if the main 
process exits before the 5 seconds
                                // have elapsed, i.e., as soon as all 
subprocesses have returned from Wait().
@@ -330,7 +326,7 @@ func launchSDKProcess() error {
                        }(pid)
                        syscall.Kill(-pid, syscall.SIGTERM)
                }
-               childPids.mu.Unlock()
+               workerMu.Unlock()
        }()
 
        var wg sync.WaitGroup
@@ -342,9 +338,9 @@ func launchSDKProcess() error {
                        bufLogger := tools.NewBufferedLogger(logger)
                        errorCount := 0
                        for {
-                               childPids.mu.Lock()
-                               if childPids.canceled {
-                                       childPids.mu.Unlock()
+                               workerMu.Lock()
+                               if shuttingDown {
+                                       workerMu.Unlock()
                                        return
                                }
 
@@ -369,8 +365,9 @@ func launchSDKProcess() error {
 
                                logger.Printf(ctx, "Executing Python (%v): %v 
%v", envStr, currentProg, strings.Join(currentArgs, " "))
                                cmd := StartCommandEnv(currentEnv, os.Stdin, 
bufLogger, bufLogger, currentProg, currentArgs...)
-                               childPids.v = append(childPids.v, 
cmd.Process.Pid)
-                               childPids.mu.Unlock()
+                               logger.Printf(ctx, "Started worker %s with PID 
%d", workerId, cmd.Process.Pid)
+                               activePids = append(activePids, cmd.Process.Pid)
+                               workerMu.Unlock()
 
                                var timer *time.Timer
                                var profilingTimedOut atomic.Bool
@@ -379,8 +376,8 @@ func launchSDKProcess() error {
                                if profilingActive && pcfg.StopAfterSec > 0 {
                                        duration := 
time.Duration(pcfg.StopAfterSec) * time.Second
                                        timer = time.AfterFunc(duration, func() 
{
-                                               childPids.mu.Lock()
-                                               defer childPids.mu.Unlock()
+                                               workerMu.Lock()
+                                               defer workerMu.Unlock()
                                                if cmd.Process != nil {
                                                        logger.Printf(ctx, 
"Profiling timeout of %d seconds reached. Sending SIGINT to worker %s",
                                                                
pcfg.StopAfterSec, workerId)
@@ -391,6 +388,7 @@ func launchSDKProcess() error {
                                }
 
                                err := cmd.Wait()
+                               unregisterPid(cmd.Process.Pid)
                                if timer != nil {
                                        timer.Stop()
                                }
@@ -417,6 +415,7 @@ func launchSDKProcess() error {
                                                logger.Warnf(ctx, "Python 
(worker %v) exited %v times: %v\nrestarting SDK process",
                                                        workerId, errorCount, 
err)
                                        } else {
+                                               cleanUpProfiler(ctx, logger)
                                                logger.Fatalf(ctx, "Python 
(worker %v) exited %v times: %v\nout of retries, failing container",
                                                        workerId, errorCount, 
err)
                                        }
@@ -595,3 +594,23 @@ func logSubmissionEnvDependencies(ctx context.Context, 
bufLogger *tools.Buffered
        bufLogger.Printf(ctx, "Dependencies in submission environment:\n%s", 
string(content))
        return nil
 }
+
+var (
+       activePids []int
+)
+
+func unregisterPid(pid int) {
+       workerMu.Lock()
+       defer workerMu.Unlock()
+       activePids = slices.DeleteFunc(activePids, func(p int) bool {
+               return p == pid
+       })
+}
+
+func getActivePids() []int {
+       workerMu.Lock()
+       defer workerMu.Unlock()
+       pids := make([]int, len(activePids))
+       copy(pids, activePids)
+       return pids
+}
diff --git a/sdks/python/container/ml/py310/base_image_requirements.txt 
b/sdks/python/container/ml/py310/base_image_requirements.txt
index bec06ec1bef..7b103880160 100644
--- a/sdks/python/container/ml/py310/base_image_requirements.txt
+++ b/sdks/python/container/ml/py310/base_image_requirements.txt
@@ -181,6 +181,7 @@ PyMySQL==1.2.0
 pyOpenSSL==26.2.0
 pyparsing==3.3.2
 pyproject_hooks==1.2.0
+pystack==1.7.0
 pytest==9.1.1
 pytest-timeout==2.4.0
 pytest-xdist==3.8.0
diff --git a/sdks/python/container/ml/py310/gpu_image_requirements.txt 
b/sdks/python/container/ml/py310/gpu_image_requirements.txt
index 4f9e02edf77..2d490ecd55d 100644
--- a/sdks/python/container/ml/py310/gpu_image_requirements.txt
+++ b/sdks/python/container/ml/py310/gpu_image_requirements.txt
@@ -256,6 +256,7 @@ PyMySQL==1.2.0
 pyOpenSSL==26.2.0
 pyparsing==3.3.2
 pyproject_hooks==1.2.0
+pystack==1.7.0
 pytest==9.1.1
 pytest-timeout==2.4.0
 pytest-xdist==3.8.0
diff --git a/sdks/python/container/ml/py311/base_image_requirements.txt 
b/sdks/python/container/ml/py311/base_image_requirements.txt
index 3f9a18099fa..58790952177 100644
--- a/sdks/python/container/ml/py311/base_image_requirements.txt
+++ b/sdks/python/container/ml/py311/base_image_requirements.txt
@@ -180,6 +180,7 @@ PyMySQL==1.2.0
 pyOpenSSL==26.2.0
 pyparsing==3.3.2
 pyproject_hooks==1.2.0
+pystack==1.7.0
 pytest==9.1.1
 pytest-timeout==2.4.0
 pytest-xdist==3.8.0
diff --git a/sdks/python/container/ml/py311/gpu_image_requirements.txt 
b/sdks/python/container/ml/py311/gpu_image_requirements.txt
index 7cb148c46f9..69e9033253c 100644
--- a/sdks/python/container/ml/py311/gpu_image_requirements.txt
+++ b/sdks/python/container/ml/py311/gpu_image_requirements.txt
@@ -255,6 +255,7 @@ PyMySQL==1.2.0
 pyOpenSSL==26.2.0
 pyparsing==3.3.2
 pyproject_hooks==1.2.0
+pystack==1.7.0
 pytest==9.1.1
 pytest-timeout==2.4.0
 pytest-xdist==3.8.0
diff --git a/sdks/python/container/ml/py312/base_image_requirements.txt 
b/sdks/python/container/ml/py312/base_image_requirements.txt
index 265696dc7e8..7b289ba98c7 100644
--- a/sdks/python/container/ml/py312/base_image_requirements.txt
+++ b/sdks/python/container/ml/py312/base_image_requirements.txt
@@ -178,6 +178,7 @@ PyMySQL==1.2.0
 pyOpenSSL==26.2.0
 pyparsing==3.3.2
 pyproject_hooks==1.2.0
+pystack==1.7.0
 pytest==9.1.1
 pytest-timeout==2.4.0
 pytest-xdist==3.8.0
diff --git a/sdks/python/container/ml/py312/gpu_image_requirements.txt 
b/sdks/python/container/ml/py312/gpu_image_requirements.txt
index 7cc83b10ca0..44444ae5fff 100644
--- a/sdks/python/container/ml/py312/gpu_image_requirements.txt
+++ b/sdks/python/container/ml/py312/gpu_image_requirements.txt
@@ -253,6 +253,7 @@ PyMySQL==1.2.0
 pyOpenSSL==26.2.0
 pyparsing==3.3.2
 pyproject_hooks==1.2.0
+pystack==1.7.0
 pytest==9.1.1
 pytest-timeout==2.4.0
 pytest-xdist==3.8.0
diff --git a/sdks/python/container/ml/py313/base_image_requirements.txt 
b/sdks/python/container/ml/py313/base_image_requirements.txt
index 67f459b4fb2..cd334be9c17 100644
--- a/sdks/python/container/ml/py313/base_image_requirements.txt
+++ b/sdks/python/container/ml/py313/base_image_requirements.txt
@@ -177,6 +177,7 @@ PyMySQL==1.2.0
 pyOpenSSL==26.2.0
 pyparsing==3.3.2
 pyproject_hooks==1.2.0
+pystack==1.7.0
 pytest==9.1.1
 pytest-timeout==2.4.0
 pytest-xdist==3.8.0
diff --git a/sdks/python/container/profiler.go 
b/sdks/python/container/profiler.go
index 64211e9fac2..d19923f912c 100644
--- a/sdks/python/container/profiler.go
+++ b/sdks/python/container/profiler.go
@@ -22,6 +22,7 @@ import (
        "os/exec"
        "path/filepath"
        "strings"
+       "sync"
        "time"
 
        "github.com/apache/beam/sdks/v2/go/container/tools"
@@ -29,7 +30,17 @@ import (
 
 type profilerConfigKeyType struct{}
 
-var profilerConfigKey profilerConfigKeyType
+var (
+       profilerConfigKey profilerConfigKeyType
+       profilerMu        sync.Mutex
+       cleanupCallbacks  []func(ctx context.Context, logger *tools.Logger)
+)
+
+// registerCleanupCallback registers a function to be executed synchronously 
during container shutdown.
+// This allows individual profiling agents to perform the final iteration of 
profile post processing.
+func registerCleanupCallback(cb func(ctx context.Context, logger 
*tools.Logger)) {
+       cleanupCallbacks = append(cleanupCallbacks, cb)
+}
 
 // ProfilerConfig holds all pre-computed profiling parameters.
 type ProfilerConfig struct {
@@ -46,6 +57,7 @@ type ProfilerConfig struct {
        StopAfterSec           int
        StopAfterCrash         bool
        PostprocessIntervalSec int
+       GcloudAvailable        bool
 }
 
 // setupProfilerConfig parses PipelineOptionsData and stores a resolved 
ProfilerConfig in the context.
@@ -73,8 +85,14 @@ func setupProfilerConfig(ctx context.Context, logger 
*tools.Logger, opts *Pipeli
        sentinelPath := filepath.Join(tempLocation, 
fmt.Sprintf(".profiler_disengaged_%s_%s", jobId, hostname))
 
        var gcsDestPath string
+       gcloudAvailable := false
        if strings.HasPrefix(opts.Options.ProfileLocation, "gs://") {
                gcsDestPath = strings.TrimSuffix(opts.Options.ProfileLocation, 
"/")
+               if _, err := exec.LookPath("gcloud"); err == nil {
+                       gcloudAvailable = true
+               } else {
+                       logger.Errorf(ctx, "gcloud is not available, profiles 
will not be uploaded.")
+               }
        }
 
        config := &ProfilerConfig{
@@ -91,6 +109,7 @@ func setupProfilerConfig(ctx context.Context, logger 
*tools.Logger, opts *Pipeli
                StopAfterSec:           opts.Options.ProfilerStopAfterSec,
                StopAfterCrash:         opts.Options.ProfilerStopAfterCrash,
                PostprocessIntervalSec: 
opts.Options.ProfilePostprocessIntervalSec,
+               GcloudAvailable:        gcloudAvailable,
        }
 
        return context.WithValue(ctx, profilerConfigKey, config)
@@ -125,29 +144,38 @@ func startProfilerBackgroundTasks(ctx context.Context, 
logger *tools.Logger) {
                logger.Warnf(ctx, "Failed to create ProfileTempLocation: %v", 
err)
        }
 
-       if pcfg.GcsDestPath != "" {
-               if _, err := exec.LookPath("gcloud"); err != nil {
-                       logger.Errorf(ctx, "gcloud is not available, profiles 
will not be uploaded.")
-               } else {
-                       if pcfg.UploadIntervalSec > 0 {
-                               go func() {
-                                       for {
-                                               select {
-                                               case <-ctx.Done():
-                                                       return
-                                               case 
<-time.After(time.Duration(pcfg.UploadIntervalSec) * time.Second):
-                                                       // TODO(tvalentyn): 
Consider a periodic cleanup as well to save local disk space.
-                                                       syncProfilesToGCS(ctx, 
logger, pcfg.BaseTempDir, pcfg.GcsDestPath)
-                                               }
+       if pcfg.GcsDestPath != "" && pcfg.GcloudAvailable {
+               if pcfg.UploadIntervalSec > 0 {
+                       go func() {
+                               for {
+                                       select {
+                                       case <-ctx.Done():
+                                               return
+                                       case 
<-time.After(time.Duration(pcfg.UploadIntervalSec) * time.Second):
+                                               // TODO(tvalentyn): Consider a 
periodic cleanup as well to save local disk space.
+                                               syncProfilesToGCS(ctx, logger, 
pcfg.BaseTempDir, pcfg.GcsDestPath)
                                        }
-                               }()
-                       }
+                               }
+                       }()
                }
        }
 
-       if pcfg.Agent == "memray" {
-               go postProcessProfilesLoop(ctx, logger, pcfg.TempLocation, 
pcfg.PostprocessIntervalSec)
+       if pcfg.PostprocessIntervalSec > 0 {
+               if pcfg.Agent == "memray" {
+                       go postProcessProfilesLoop(ctx, logger, pcfg)
+                       registerCleanupCallback(func(ctx context.Context, 
logger *tools.Logger) {
+                               runPostProcessingSweep(ctx, logger, 
pcfg.TempLocation, pcfg.PostprocessIntervalSec)
+                       })
+               }
+
+               if pcfg.Agent == "coredump" {
+                       go monitorCoredumpsLoop(ctx, logger, pcfg)
+                       registerCleanupCallback(func(ctx context.Context, 
logger *tools.Logger) {
+                               processNewCoredumps(ctx, logger, pcfg)
+                       })
+               }
        }
+
 }
 
 // maybeWithProfiler builds the execution arguments and environment variables 
if profiling is enabled and active.
@@ -192,6 +220,9 @@ func maybeWithProfiler(
                }
                env["HEAPPROFILE"] = tcmallocHeapPath
                args = currentArgs
+       } else if pcfg.Agent == "coredump" {
+               // No wrapping of the executable is needed for coredump 
analysis.
+               args = currentArgs
        } else {
                prog = pcfg.Agent
                args = append(append([]string{}, pcfg.ExtraArgs...), 
currentProg)
@@ -223,6 +254,14 @@ func stopProfiling(ctx context.Context) error {
        return err
 }
 
+// isProfilerDisengaged checks if the stop sentinel file exists.
+func isProfilerDisengaged(pcfg *ProfilerConfig) bool {
+       if _, err := os.Stat(pcfg.StopSentinelPath); err == nil {
+               return true
+       }
+       return false
+}
+
 // syncProfilesToGCS uploads newly created local memory profiles to the 
designated GCS target path using gcloud storage.
 func syncProfilesToGCS(ctx context.Context, logger *tools.Logger, localDir, 
gcsDest string) {
        entries, err := os.ReadDir(localDir)
@@ -241,18 +280,18 @@ func syncProfilesToGCS(ctx context.Context, logger 
*tools.Logger, localDir, gcsD
 }
 
 // postProcessProfilesLoop runs a background loop that periodically triggers 
profile post-processing if enabled.
-func postProcessProfilesLoop(ctx context.Context, logger *tools.Logger, 
profilesDir string, intervalSec int) {
-       if intervalSec <= 0 {
-               return
-       }
-
+func postProcessProfilesLoop(ctx context.Context, logger *tools.Logger, pcfg 
*ProfilerConfig) {
        for {
-               runPostProcessingSweep(ctx, logger, profilesDir, intervalSec)
+               runPostProcessingSweep(ctx, logger, pcfg.TempLocation, 
pcfg.PostprocessIntervalSec)
+
+               if isProfilerDisengaged(pcfg) {
+                       return
+               }
 
                select {
                case <-ctx.Done():
                        return
-               case <-time.After(time.Duration(intervalSec) * time.Second):
+               case <-time.After(time.Duration(pcfg.PostprocessIntervalSec) * 
time.Second):
                        // Block until the sleep completes before starting the 
next sweep
                }
        }
@@ -260,6 +299,9 @@ func postProcessProfilesLoop(ctx context.Context, logger 
*tools.Logger, profiles
 
 // runPostProcessingSweep scans the profiles directory and launches sequential 
postprocessing for newly updated profiles.
 func runPostProcessingSweep(ctx context.Context, logger *tools.Logger, 
profilesDir string, intervalSec int) {
+       profilerMu.Lock()
+       defer profilerMu.Unlock()
+
        files, err := os.ReadDir(profilesDir)
        if err != nil {
                return
@@ -334,3 +376,212 @@ func needsProcessing(binInfo os.FileInfo, path string) 
bool {
        // Don't regenerate when there were no updates to the profile.
        return binInfo.ModTime().After(info.ModTime())
 }
+
+func monitorCoredumpsLoop(ctx context.Context, logger *tools.Logger, pcfg 
*ProfilerConfig) {
+       if pcfg.PostprocessIntervalSec <= 0 {
+               return
+       }
+
+       interval := time.Duration(pcfg.PostprocessIntervalSec) * time.Second
+       logger.Printf(ctx, "Monitoring core dumps every %v", interval)
+
+       ticker := time.NewTicker(interval)
+       defer ticker.Stop()
+
+       for {
+               select {
+               case <-ctx.Done():
+                       return
+               case <-ticker.C:
+                       processNewCoredumps(ctx, logger, pcfg)
+                       if isProfilerDisengaged(pcfg) {
+                               return
+                       }
+               }
+       }
+}
+
+func processNewCoredumps(ctx context.Context, logger *tools.Logger, pcfg 
*ProfilerConfig) {
+       profilerMu.Lock()
+       defer profilerMu.Unlock()
+
+       // We expect the runner runtime environment to set the core pattern
+       // to /tmp/beam_coredump.%e.%p or similar. To do that, we pass
+       // the --experiment=core_pattern pipeline option, which can be 
interpreted by a runner.
+       coreDir := "/tmp"
+       files, err := os.ReadDir(coreDir)
+       if err != nil {
+               return
+       }
+
+       prefix := "beam_coredump."
+
+       for _, file := range files {
+               if file.IsDir() {
+                       continue
+               }
+               name := file.Name()
+               if !strings.HasPrefix(name, prefix) {
+                       continue
+               }
+
+               corePath := filepath.Join(coreDir, name)
+               var info os.FileInfo
+               var err error
+
+               for {
+                       info, err = os.Stat(corePath)
+                       if err != nil || time.Since(info.ModTime()) >= 
2*time.Second {
+                               break
+                       }
+                       // Wait for the core file to finish being written.
+                       time.Sleep(500 * time.Millisecond)
+               }
+               if err != nil {
+                       continue
+               }
+
+               logger.Printf(ctx, "Found core dump file: %s (%d bytes)", name, 
info.Size())
+
+               // Find python executable. Since the worker might be running in 
a venv,
+               // we look for "python" in the PATH.
+               pythonProg := "python"
+               if path, err := exec.LookPath("python"); err == nil {
+                       pythonProg = path
+               }
+
+               timeSuffix := info.ModTime().Format("20060102150405")
+               newName := fmt.Sprintf("%s-%s", name, timeSuffix)
+               destTxtPath := filepath.Join(pcfg.TempLocation, 
fmt.Sprintf("%s.txt", newName))
+
+               // Delete the core file after up to 2 attempts to process it.
+               shouldDelete := time.Since(info.ModTime()) > 
time.Duration(pcfg.PostprocessIntervalSec)*time.Second
+
+               pystackPath, pystackErr := exec.LookPath("pystack")
+               gdbPath, gdbErr := exec.LookPath("gdb")
+
+               if pystackErr != nil && gdbErr != nil {
+                       logger.Warnf(ctx, "Core dump analysis enabled but no 
analysis tools found. Please install pystack (recommended) or/and gdb into the 
runtime environment.")
+               }
+
+               if pystackErr == nil {
+                       args := []string{"core"}
+                       if len(pcfg.ExtraArgs) > 0 {
+                               args = append(args, pcfg.ExtraArgs...)
+                       } else {
+                               args = append(args, "--native-last")
+                       }
+                       args = append(args, corePath, pythonProg)
+
+                       logger.Printf(ctx, "Running pystack %s", 
strings.Join(args, " "))
+                       cmd := exec.CommandContext(ctx, pystackPath, args...)
+                       output, err := cmd.CombinedOutput()
+                       if err != nil {
+                               logger.Warnf(ctx, "pystack failed on %s: %v. 
Output:\n%s", name, err, string(output))
+                       } else {
+                               if err := os.WriteFile(destTxtPath, output, 
0644); err != nil {
+                                       logger.Warnf(ctx, "Failed to write 
pystack output to %s: %v", destTxtPath, err)
+                               }
+                               pystackSummary := 
createPystackSummary(string(output))
+                               logger.Errorf(ctx, "Full pystack coredump 
analysis saved to %s.txt\nExcerpt:\n%s", newName, pystackSummary)
+                               shouldDelete = true
+                       }
+               }
+
+               if gdbErr == nil {
+                       gdbArgs := []string{
+                               "-batch",
+                               "-ex", "set pagination off",
+                               "-ex", "set trace-commands on",
+                               "-ex", "info sharedlibrary",
+                               "-ex", "info proc mappings",
+                               "-ex", "info threads",
+                               "-ex", "thread",
+                               "-ex", "print $_siginfo",
+                               "-ex", "info registers",
+                               "-ex", "x/10i $pc",
+                               "-ex", "x/16gx $rsp",
+                               "-ex", "bt full",
+                               "-ex", "thread apply all bt full",
+                               pythonProg,
+                               corePath,
+                       }
+                       logger.Printf(ctx, "Running gdb on %s using %s", name, 
pythonProg)
+                       gdbCmd := exec.CommandContext(ctx, gdbPath, gdbArgs...)
+                       gdbOutput, err := gdbCmd.CombinedOutput()
+                       destGdbPath := filepath.Join(pcfg.TempLocation, 
fmt.Sprintf("%s.gdb.txt", newName))
+                       if err != nil {
+                               logger.Warnf(ctx, "gdb failed on %s: %v. 
Output:\n%s", name, err, string(gdbOutput))
+                       } else {
+                               if err := os.WriteFile(destGdbPath, gdbOutput, 
0644); err != nil {
+                                       logger.Warnf(ctx, "Failed to write gdb 
output to %s: %v", destGdbPath, err)
+                               }
+                               logger.Errorf(ctx, "Full GDB coredump analysis 
saved to %s.gdb.txt", newName)
+                               shouldDelete = true
+                       }
+               }
+
+               if shouldDelete {
+                       if err := os.Remove(corePath); err != nil {
+                               logger.Warnf(ctx, "Failed to delete core dump 
%s: %v", corePath, err)
+                       }
+               }
+       }
+}
+
+func extractGILThread(output string) string {
+       lines := strings.Split(output, "\n")
+       var result []string
+       recording := false
+       for _, line := range lines {
+               trimmed := strings.TrimSpace(line)
+               if strings.Contains(line, "Has the GIL") {
+                       recording = true
+               }
+               if recording {
+                       result = append(result, line)
+                       if trimmed == "" {
+                               break
+                       }
+               }
+       }
+       if len(result) == 0 {
+               return ""
+       }
+       return strings.Join(result, "\n")
+}
+
+func firstNLines(s string, n int) string {
+       lines := strings.Split(s, "\n")
+       if len(lines) <= n {
+               return s
+       }
+       return strings.Join(lines[:n], "\n")
+}
+
+func createPystackSummary(output string) string {
+       gilThreadTrace := extractGILThread(output)
+       if gilThreadTrace != "" {
+               return gilThreadTrace
+       }
+       return firstNLines(output, 100)
+}
+
+// cleanUpProfiler checks for and uploads any final profiler artifacts before 
container exit.
+func cleanUpProfiler(ctx context.Context, logger *tools.Logger) {
+       pcfg := getProfilerConfig(ctx)
+       if pcfg == nil || !pcfg.Enabled {
+               return
+       }
+
+       logger.Printf(ctx, "Running final profiler cleanup sweep and GCS 
sync...")
+
+       // Execute all registered agent-specific cleanups
+       for _, cb := range cleanupCallbacks {
+               cb(ctx, logger)
+       }
+
+       if pcfg.GcsDestPath != "" && pcfg.GcloudAvailable {
+               syncProfilesToGCS(ctx, logger, pcfg.BaseTempDir, 
pcfg.GcsDestPath)
+       }
+}
diff --git a/sdks/python/container/profiler_test.go 
b/sdks/python/container/profiler_test.go
new file mode 100644
index 00000000000..27abf8a2ab3
--- /dev/null
+++ b/sdks/python/container/profiler_test.go
@@ -0,0 +1,125 @@
+// Licensed to the Apache Software Foundation (ASF) under one or more
+// contributor license agreements.  See the NOTICE file distributed with
+// this work for additional information regarding copyright ownership.
+// The ASF licenses this file to You under the Apache License, Version 2.0
+// (the "License"); you may not use this file except in compliance with
+// the License.  You may obtain a copy of the License at
+//
+//    http://www.apache.org/licenses/LICENSE-2.0
+//
+// Unless required by applicable law or agreed to in writing, software
+// distributed under the License is distributed on an "AS IS" BASIS,
+// WITHOUT WARRANTIES OR CONDITIONS OF ANY KIND, either express or implied.
+// See the License for the specific language governing permissions and
+// limitations under the License.
+
+package main
+
+import (
+       "context"
+       "os"
+       "path/filepath"
+       "testing"
+)
+
+func TestActivePidsRegistry(t *testing.T) {
+       // Reset active PIDs
+       activePids = nil
+
+       activePids = append(activePids, 101)
+       activePids = append(activePids, 102)
+
+       pids := getActivePids()
+       if len(pids) != 2 || pids[0] != 101 || pids[1] != 102 {
+               t.Errorf("Expected active pids [101, 102], got %v", pids)
+       }
+
+       unregisterPid(101)
+       pids = getActivePids()
+       if len(pids) != 1 || pids[0] != 102 {
+               t.Errorf("Expected active pids [102], got %v", pids)
+       }
+
+       unregisterPid(102)
+       pids = getActivePids()
+       if len(pids) != 0 {
+               t.Errorf("Expected active pids empty, got %v", pids)
+       }
+}
+
+func TestSetupProfilerConfig(t *testing.T) {
+       opts := &PipelineOptionsData{
+               Options: OptionsData{
+                       ProfilerAgent: "coredump",
+                       JobId:         "test-job",
+               },
+       }
+       ctx := setupProfilerConfig(context.Background(), nil, opts)
+       pcfg := getProfilerConfig(ctx)
+       if pcfg == nil {
+               t.Fatal("ProfilerConfig was nil")
+       }
+
+       if pcfg.Agent != "coredump" {
+               t.Errorf("Expected agent coredump, got %s", pcfg.Agent)
+       }
+}
+
+func TestIsProfilerDisengaged(t *testing.T) {
+       tempDir, err := os.MkdirTemp("", "disengage_test")
+       if err != nil {
+               t.Fatal(err)
+       }
+       defer os.RemoveAll(tempDir)
+
+       sentinelPath := filepath.Join(tempDir, "stop_sentinel")
+       pcfg := &ProfilerConfig{
+               StopSentinelPath: sentinelPath,
+       }
+
+       if isProfilerDisengaged(pcfg) {
+               t.Error("Expected profiler NOT to be disengaged before sentinel 
creation")
+       }
+
+       // Create sentinel file
+       if err := os.WriteFile(sentinelPath, []byte{}, 0644); err != nil {
+               t.Fatal(err)
+       }
+
+       if !isProfilerDisengaged(pcfg) {
+               t.Error("Expected profiler to be disengaged after sentinel 
creation")
+       }
+}
+
+func TestCreatePystackSummary(t *testing.T) {
+       t.Run("ExtractsGILThreadTrace", func(t *testing.T) {
+               output := "Thread 1 (waiting):\n" +
+                       "  File \"worker.py\", line 10, in run\n" +
+                       "\n" +
+                       "Thread 2 (active, Has the GIL):\n" +
+                       "  File \"main.py\", line 42, in execute\n" +
+                       "  File \"db.py\", line 5, in query\n" +
+                       "\n" +
+                       "Thread 3 (idle):\n" +
+                       "  File \"server.py\", line 99, in listen\n"
+
+               expected := "Thread 2 (active, Has the GIL):\n" +
+                       "  File \"main.py\", line 42, in execute\n" +
+                       "  File \"db.py\", line 5, in query\n"
+
+               result := createPystackSummary(output)
+               if result != expected {
+                       t.Errorf("Expected:\n%s\nGot:\n%s", expected, result)
+               }
+       })
+
+       t.Run("FallbackSmallOutput", func(t *testing.T) {
+               output := "Thread 1 (waiting):\n" +
+                       "  File \"worker.py\", line 10, in run"
+
+               result := createPystackSummary(output)
+               if result != output {
+                       t.Errorf("Expected identical output, got:\n%s", result)
+               }
+       })
+}
diff --git a/sdks/python/container/py310/base_image_requirements.txt 
b/sdks/python/container/py310/base_image_requirements.txt
index 85e7615cfdb..dc78e342a91 100644
--- a/sdks/python/container/py310/base_image_requirements.txt
+++ b/sdks/python/container/py310/base_image_requirements.txt
@@ -163,6 +163,7 @@ PyMySQL==1.2.0
 pyOpenSSL==26.2.0
 pyparsing==3.3.2
 pyproject_hooks==1.2.0
+pystack==1.7.0
 pytest==9.1.1
 pytest-timeout==2.4.0
 pytest-xdist==3.8.0
diff --git a/sdks/python/container/py311/base_image_requirements.txt 
b/sdks/python/container/py311/base_image_requirements.txt
index b9ede50e41d..61ff31bbd4b 100644
--- a/sdks/python/container/py311/base_image_requirements.txt
+++ b/sdks/python/container/py311/base_image_requirements.txt
@@ -162,6 +162,7 @@ PyMySQL==1.2.0
 pyOpenSSL==26.2.0
 pyparsing==3.3.2
 pyproject_hooks==1.2.0
+pystack==1.7.0
 pytest==9.1.1
 pytest-timeout==2.4.0
 pytest-xdist==3.8.0
diff --git a/sdks/python/container/py312/base_image_requirements.txt 
b/sdks/python/container/py312/base_image_requirements.txt
index 4edcb642110..cd93b7fc3c0 100644
--- a/sdks/python/container/py312/base_image_requirements.txt
+++ b/sdks/python/container/py312/base_image_requirements.txt
@@ -160,6 +160,7 @@ PyMySQL==1.2.0
 pyOpenSSL==26.2.0
 pyparsing==3.3.2
 pyproject_hooks==1.2.0
+pystack==1.7.0
 pytest==9.1.1
 pytest-timeout==2.4.0
 pytest-xdist==3.8.0
diff --git a/sdks/python/container/py313/base_image_requirements.txt 
b/sdks/python/container/py313/base_image_requirements.txt
index a9728cd5e10..555de648d6b 100644
--- a/sdks/python/container/py313/base_image_requirements.txt
+++ b/sdks/python/container/py313/base_image_requirements.txt
@@ -159,6 +159,7 @@ PyMySQL==1.2.0
 pyOpenSSL==26.2.0
 pyparsing==3.3.2
 pyproject_hooks==1.2.0
+pystack==1.7.0
 pytest==9.1.1
 pytest-timeout==2.4.0
 pytest-xdist==3.8.0
diff --git a/sdks/python/container/py314/base_image_requirements.txt 
b/sdks/python/container/py314/base_image_requirements.txt
index 1598d940e93..bf74a89916c 100644
--- a/sdks/python/container/py314/base_image_requirements.txt
+++ b/sdks/python/container/py314/base_image_requirements.txt
@@ -158,6 +158,7 @@ PyMySQL==1.2.0
 pyOpenSSL==26.2.0
 pyparsing==3.3.2
 pyproject_hooks==1.2.0
+pystack==1.7.0
 pytest==9.1.1
 pytest-timeout==2.4.0
 pytest-xdist==3.8.0

Reply via email to