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 {
}