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

jrmccluskey 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 be4211d1d24 [Go SDK] Recreate a failed data/state channel without 
holding the lock (#40415)
be4211d1d24 is described below

commit be4211d1d2408e1571b752a884052f301781793a
Author: Igor <[email protected]>
AuthorDate: Wed Oct 7 18:34:10 2026 +0300

    [Go SDK] Recreate a failed data/state channel without holding the lock 
(#40415)
    
    * [Go SDK] Recreate a failed data/state channel without holding the manager 
lock
    
    * [Go SDK] Add CHANGES for #40414 and fix SA2001 in recreate test
    
    * [Go SDK] Keep StateChannel recreate on the manager lock
    
    * [Go SDK] Test data channel recreate with production Open
    
    * [Go SDK] Release the manager lock before closeInstruction
---
 CHANGES.md                                         |   1 +
 sdks/go/pkg/beam/core/runtime/harness/datamgr.go   |  13 +-
 .../pkg/beam/core/runtime/harness/datamgr_test.go  | 188 +++++++++++++++++++++
 3 files changed, 197 insertions(+), 5 deletions(-)

diff --git a/CHANGES.md b/CHANGES.md
index 8aa31f5bfb7..ccf5b7f142d 100644
--- a/CHANGES.md
+++ b/CHANGES.md
@@ -85,6 +85,7 @@
 
 * (Go) Fixed a data race on the Prism runner's artifact cache map in 
JobServices ([#32656](https://github.com/apache/beam/issues/32656)).
 * (Go) Fixed the harness leaking Data/State gRPC streams after the worker 
stops, and a deadlock when Send returns EOF 
([#40260](https://github.com/apache/beam/issues/40260)).
+* (Go) Fixed a deadlock recreating a data channel while holding the channel 
lock ([#40414](https://github.com/apache/beam/issues/40414)).
 * (Java) Fixed the declared schema of the error output of the Kafka write 
SchemaTransform, which wrapped the error schema a second time and did not match 
the rows it emits ([#39760](https://github.com/apache/beam/issues/39760)).
 * (Go) Fixed pubsubio importing a `google.golang.org/genproto` package removed 
in recent releases, which broke builds of Go modules depending on a current 
`genproto` version ([#40018](https://github.com/apache/beam/issues/40018)).
 * (Java) BigQueryIO now treats a 404 when deleting a temporary table or 
dataset as success, so a replayed work item whose earlier attempt already 
deleted it no longer retries forever 
([#24997](https://github.com/apache/beam/issues/24997)).
diff --git a/sdks/go/pkg/beam/core/runtime/harness/datamgr.go 
b/sdks/go/pkg/beam/core/runtime/harness/datamgr.go
index b9718e9b42b..66ebd5aa6f9 100644
--- a/sdks/go/pkg/beam/core/runtime/harness/datamgr.go
+++ b/sdks/go/pkg/beam/core/runtime/harness/datamgr.go
@@ -159,14 +159,17 @@ func (m *DataChannelManager) Close() {
 }
 
 func (m *DataChannelManager) closeInstruction(instID instructionID, ports 
[]exec.Port) error {
+       // copy channels so removeInstruction can take ch.mu without holding 
m.mu.
+       var chans []*DataChannel
        m.mu.Lock()
-       defer m.mu.Unlock()
-       var firstNonNilError error
        for _, port := range ports {
-               ch, ok := m.ports[port.URL]
-               if !ok {
-                       continue
+               if ch, ok := m.ports[port.URL]; ok {
+                       chans = append(chans, ch)
                }
+       }
+       m.mu.Unlock()
+       var firstNonNilError error
+       for _, ch := range chans {
                err := ch.removeInstruction(instID)
                if err != nil && firstNonNilError == nil {
                        firstNonNilError = err
diff --git a/sdks/go/pkg/beam/core/runtime/harness/datamgr_test.go 
b/sdks/go/pkg/beam/core/runtime/harness/datamgr_test.go
index 637076d6e44..b48967daf9c 100644
--- a/sdks/go/pkg/beam/core/runtime/harness/datamgr_test.go
+++ b/sdks/go/pkg/beam/core/runtime/harness/datamgr_test.go
@@ -21,6 +21,7 @@ import (
        "fmt"
        "io"
        "log"
+       "net"
        "runtime"
        "strings"
        "sync"
@@ -30,6 +31,7 @@ import (
 
        "github.com/apache/beam/sdks/v2/go/pkg/beam/core/runtime/exec"
        fnpb "github.com/apache/beam/sdks/v2/go/pkg/beam/model/fnexecution_v1"
+       "google.golang.org/grpc"
 )
 
 const extraData = 2
@@ -644,6 +646,192 @@ func TestTimerWriterSendEOF(t *testing.T) {
        }
 }
 
+// Production Open against closeInstruction.
+// Flush holds ch.mu inside client.Send.
+// Close takes m.mu and waits on ch.mu.
+// Then the stream fails.
+func TestDataChannelTerminate_recreate(t *testing.T) {
+       hs := newHoldDataServer()
+       lis := newStallListener(t)
+       gs := grpc.NewServer()
+       fnpb.RegisterBeamFnDataServer(gs, hs)
+       go gs.Serve(lis)
+       defer func() {
+               lis.fail()
+               gs.Stop()
+       }()
+
+       ctx, cancel := context.WithCancel(context.Background())
+       defer cancel()
+       m := &DataChannelManager{}
+       s := NewScopedDataManager(m, "inst1")
+       w, err := s.OpenWrite(ctx, exec.StreamID{Port: exec.Port{URL: 
lis.Addr().String()}, PtransformID: "pt"})
+       if err != nil {
+               t.Fatal(err)
+       }
+       <-hs.entered
+       lis.stall()
+
+       flushDone := make(chan error, 1)
+       go func() {
+               buf := make([]byte, 4e6)
+               for {
+                       if _, err := w.Write(buf); err != nil {
+                               flushDone <- err
+                               return
+                       }
+                       if _, err := w.Write([]byte{1}); err != nil {
+                               flushDone <- err
+                               return
+                       }
+               }
+       }()
+       select {
+       case err := <-flushDone:
+               t.Fatalf("flush returned before Close: %v", err)
+       case <-time.After(500 * time.Millisecond):
+       }
+
+       closeDone := make(chan error, 1)
+       go func() { closeDone <- s.Close() }()
+       time.Sleep(200 * time.Millisecond)
+       lis.fail()
+
+       timeout := time.After(3 * time.Second)
+       gotFlush, gotClose := false, false
+       for !gotFlush || !gotClose {
+               select {
+               case <-flushDone:
+                       gotFlush = true
+               case <-closeDone:
+                       gotClose = true
+               case <-timeout:
+                       t.Fatal("recreate and closeInstruction deadlocked")
+               }
+       }
+}
+
+// After a stream fails, the next Open must dial a new channel.
+// It must not return the channel whose stream is already dead.
+func TestDataChannelOpen_skipsFailedChannel(t *testing.T) {
+       hs := newHoldDataServer()
+       lis, err := net.Listen("tcp", "127.0.0.1:0")
+       if err != nil {
+               t.Fatal(err)
+       }
+       gs := grpc.NewServer()
+       fnpb.RegisterBeamFnDataServer(gs, hs)
+       go gs.Serve(lis)
+       defer gs.Stop()
+
+       ctx, cancel := context.WithCancel(context.Background())
+       defer cancel()
+       m := &DataChannelManager{}
+       port := exec.Port{URL: lis.Addr().String()}
+       ch1, err := m.Open(ctx, port)
+       if err != nil {
+               t.Fatal(err)
+       }
+       <-hs.entered
+       ch1.mu.Lock()
+       ch1.terminateStreamOnError(io.EOF)
+       ch1.mu.Unlock()
+
+       ch2, err := m.Open(ctx, port)
+       if err != nil {
+               t.Fatal(err)
+       }
+       if ch2 == ch1 {
+               t.Fatal("Open returned a failed channel")
+       }
+}
+
+type holdDataServer struct {
+       fnpb.UnimplementedBeamFnDataServer
+       entered chan struct{}
+       once    sync.Once
+}
+
+func newHoldDataServer() *holdDataServer {
+       return &holdDataServer{entered: make(chan struct{})}
+}
+
+func (s *holdDataServer) Data(stream fnpb.BeamFnData_DataServer) error {
+       s.once.Do(func() { close(s.entered) })
+       <-stream.Context().Done()
+       return stream.Context().Err()
+}
+
+type stallListener struct {
+       net.Listener
+       mu    sync.Mutex
+       conns []*stallConn
+}
+
+func newStallListener(t *testing.T) *stallListener {
+       t.Helper()
+       l, err := net.Listen("tcp", "127.0.0.1:0")
+       if err != nil {
+               t.Fatal(err)
+       }
+       return &stallListener{Listener: l}
+}
+
+func (l *stallListener) Accept() (net.Conn, error) {
+       c, err := l.Listener.Accept()
+       if err != nil {
+               return nil, err
+       }
+       if tc, ok := c.(*net.TCPConn); ok {
+               _ = tc.SetReadBuffer(2048)
+               _ = tc.SetWriteBuffer(2048)
+       }
+       sc := &stallConn{Conn: c, closed: make(chan struct{})}
+       l.mu.Lock()
+       l.conns = append(l.conns, sc)
+       l.mu.Unlock()
+       return sc, nil
+}
+
+func (l *stallListener) stall() {
+       l.mu.Lock()
+       defer l.mu.Unlock()
+       for _, c := range l.conns {
+               c.stalled.Store(true)
+       }
+}
+
+func (l *stallListener) fail() {
+       l.mu.Lock()
+       defer l.mu.Unlock()
+       for _, c := range l.conns {
+               c.fail()
+       }
+}
+
+type stallConn struct {
+       net.Conn
+       stalled atomic.Bool
+       once    sync.Once
+       closed  chan struct{}
+}
+
+func (c *stallConn) Read(p []byte) (int, error) {
+       if c.stalled.Load() {
+               <-c.closed
+               return 0, io.EOF
+       }
+       return c.Conn.Read(p)
+}
+
+func (c *stallConn) fail() {
+       c.stalled.Store(true)
+       c.once.Do(func() {
+               _ = c.Conn.Close()
+               close(c.closed)
+       })
+}
+
 type noopDataClient struct {
 }
 

Reply via email to