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

lostluck 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 db063b8de70 Fix data race on Prism runner's artifact cache map (#40197)
db063b8de70 is described below

commit db063b8de70e81a0080d3421beed9022270c877f
Author: waneal <[email protected]>
AuthorDate: Tue Sep 22 00:17:59 2026 +0900

    Fix data race on Prism runner's artifact cache map (#40197)
    
    s.artifacts was read/written without a lock from
    ReverseArtifactRetrievalService and GetArtifact, causing "concurrent
    map writes" under the race detector. Guard it with a dedicated
    sync.RWMutex.
    
    Fixes #32656
---
 CHANGES.md                                                   |  1 +
 .../pkg/beam/runners/prism/internal/jobservices/artifact.go  |  7 ++++---
 .../go/pkg/beam/runners/prism/internal/jobservices/server.go | 12 +++++++-----
 3 files changed, 12 insertions(+), 8 deletions(-)

diff --git a/CHANGES.md b/CHANGES.md
index 5480ddc1c21..4866998bb05 100644
--- a/CHANGES.md
+++ b/CHANGES.md
@@ -80,6 +80,7 @@
 
 ## Bugfixes
 
+* (Go) Fixed a data race on the Prism runner's artifact cache map in 
JobServices ([#32656](https://github.com/apache/beam/issues/32656)).
 * Fixed X (Java/Python) ([#X](https://github.com/apache/beam/issues/X)).
 
 ## Security Fixes
diff --git a/sdks/go/pkg/beam/runners/prism/internal/jobservices/artifact.go 
b/sdks/go/pkg/beam/runners/prism/internal/jobservices/artifact.go
index f3e97406027..babbcbf562b 100644
--- a/sdks/go/pkg/beam/runners/prism/internal/jobservices/artifact.go
+++ b/sdks/go/pkg/beam/runners/prism/internal/jobservices/artifact.go
@@ -83,10 +83,9 @@ func (s *Server) ReverseArtifactRetrievalService(stream 
jobpb.ArtifactStagingSer
                                        return err
                                }
                        }
-                       if len(s.artifacts) == 0 {
-                               s.artifacts = map[string][]byte{}
-                       }
+                       s.artifactsMu.Lock()
                        s.artifacts[string(dep.GetTypePayload())] = buf.Bytes()
+                       s.artifactsMu.Unlock()
                }
        }
        return nil
@@ -100,7 +99,9 @@ func (s *Server) ResolveArtifacts(_ context.Context, req 
*jobpb.ResolveArtifacts
 
 func (s *Server) GetArtifact(req *jobpb.GetArtifactRequest, stream 
jobpb.ArtifactRetrievalService_GetArtifactServer) error {
        info := req.GetArtifact()
+       s.artifactsMu.RLock()
        buf, ok := s.artifacts[string(info.GetTypePayload())]
+       s.artifactsMu.RUnlock()
        if !ok {
                pt := prototext.Format(info)
                slog.Warn("unable to provide artifact to worker", 
"artifact_info", pt)
diff --git a/sdks/go/pkg/beam/runners/prism/internal/jobservices/server.go 
b/sdks/go/pkg/beam/runners/prism/internal/jobservices/server.go
index 2399fd726da..f6c88205146 100644
--- a/sdks/go/pkg/beam/runners/prism/internal/jobservices/server.go
+++ b/sdks/go/pkg/beam/runners/prism/internal/jobservices/server.go
@@ -61,7 +61,8 @@ type Server struct {
        execute func(*Job)
 
        // Artifact hack
-       artifacts map[string][]byte
+       artifacts   map[string][]byte
+       artifactsMu sync.RWMutex
 
        mw *worker.MultiplexW
 }
@@ -73,10 +74,11 @@ func NewServer(port int, execute func(*Job)) *Server {
                panic(fmt.Sprintf("failed to listen: %v", err))
        }
        s := &Server{
-               lis:     lis,
-               jobs:    make(map[string]*Job),
-               execute: execute,
-               logger:  slog.Default(), // TODO substitute with a configured 
logger.
+               lis:       lis,
+               jobs:      make(map[string]*Job),
+               execute:   execute,
+               artifacts: map[string][]byte{},
+               logger:    slog.Default(), // TODO substitute with a configured 
logger.
        }
        s.logger.Info("Serving JobManagement", slog.String("endpoint", 
s.Endpoint()))
        opts := []grpc.ServerOption{

Reply via email to