jrmccluskey commented on code in PR #40415:
URL: https://github.com/apache/beam/pull/40415#discussion_r4197267113


##########
sdks/go/pkg/beam/core/runtime/harness/datamgr_test.go:
##########
@@ -644,6 +645,74 @@ func TestTimerWriterSendEOF(t *testing.T) {
        }
 }
 
+func TestDataChannelTerminate_recreate(t *testing.T) {

Review Comment:
   The unit test cases here pass on master as-is, so we're not really covering 
the new behavior at all here _or_ adding something that functions as a 
regression test for the reported bug. These are fine to include, but a repro 
unit test would be much more useful. 



##########
sdks/go/pkg/beam/core/runtime/harness/statemgr.go:
##########
@@ -619,9 +619,15 @@ func (m *StateChannelManager) Open(ctx context.Context, 
port exec.Port) (*StateC
                default:
                        log.Warnf(ctx, "forcing StateChannel[%v] reconnection 
on port %v due to %v", id, port, err)
                }
-               m.mu.Lock()
-               delete(m.ports, port.URL)
-               m.mu.Unlock()
+               // Remove this channel from the port map after releasing the 
channel lock.
+               // Keep the mapping if Open has already stored a replacement.
+               go func() {
+                       m.mu.Lock()
+                       if m.ports[port.URL] == ch {
+                               delete(m.ports, port.URL)
+                       }
+                       m.mu.Unlock()
+               }()

Review Comment:
   I don't think the StateChannelManager encounters the same problem here, 
since there isn't an equivalent of closeInstruction (nothing holds m.mu and 
blocks waiting for ch.mu here, Close() releases m.mu before acquiring any 
channel lock.) If you have a clear repro here we can revisit it



##########
sdks/go/pkg/beam/core/runtime/harness/statemgr_test.go:
##########
@@ -585,6 +585,82 @@ func TestStateChannelWriteEOF(t *testing.T) {
        }
 }
 
+func TestStateChannel_recreate(t *testing.T) {

Review Comment:
   Same thing with this slate of unit tests, we're not covering any sort of 
regression behavior



##########
sdks/go/pkg/beam/core/runtime/harness/datamgr.go:
##########
@@ -136,9 +136,15 @@ func (m *DataChannelManager) Open(ctx context.Context, 
port exec.Port) (*DataCha
                default:
                        log.Warnf(ctx, "forcing DataChannel[%v] reconnection on 
port %v due to %v", id, port, err)
                }
-               m.mu.Lock()
-               delete(m.ports, port.URL)
-               m.mu.Unlock()
+               // Remove this channel from the port map after releasing ch.mu.
+               // Keep the mapping if Open has already stored a replacement.
+               go func() {
+                       m.mu.Lock()
+                       if m.ports[port.URL] == ch {
+                               delete(m.ports, port.URL)
+                       }
+                       m.mu.Unlock()
+               }()

Review Comment:
   Spinning up a goroutine here might actually introduce a different race 
condition, since there's no guarantee that the other thread will execute right 
away. If Open() is called on the port before deleting the dead channel, it will 
return the dead channel and usage downstream will fail



-- 
This is an automated message from the Apache Git Service.
To respond to the message, please log on to GitHub and use the
URL above to go to the specific comment.

To unsubscribe, e-mail: [email protected]

For queries about this service, please contact Infrastructure at:
[email protected]

Reply via email to