@@ -78,7 +78,7 @@ final class StandardIO: ManagedProcess.IO & Sendable {
7878 port: stdinPort,
7979 cid: VsockType . hostCID
8080 )
81- let stdinSocket = try Socket ( type: type)
81+ let stdinSocket = try Socket ( type: type, closeOnDeinit : false )
8282 try stdinSocket. connect ( )
8383 self . stdinSocket = stdinSocket
8484
@@ -93,7 +93,8 @@ final class StandardIO: ManagedProcess.IO & Sendable {
9393 port: stdoutPort,
9494 cid: VsockType . hostCID
9595 )
96- let stdoutSocket = try Socket ( type: type)
96+ // These fd's get closed when cleanupRelay is called
97+ let stdoutSocket = try Socket ( type: type, closeOnDeinit: false )
9798 try stdoutSocket. connect ( )
9899 self . stdoutSocket = stdoutSocket
99100
@@ -108,7 +109,7 @@ final class StandardIO: ManagedProcess.IO & Sendable {
108109 port: stderrPort,
109110 cid: VsockType . hostCID
110111 )
111- let stderrSocket = try Socket ( type: type)
112+ let stderrSocket = try Socket ( type: type, closeOnDeinit : false )
112113 try stderrSocket. connect ( )
113114 self . stderrSocket = stderrSocket
114115
@@ -125,12 +126,19 @@ final class StandardIO: ManagedProcess.IO & Sendable {
125126 func relay( readFromFd: Int32 , writeToFd: Int32 ) throws {
126127 let readFrom = OSFile ( fd: readFromFd)
127128 let writeTo = OSFile ( fd: writeToFd)
128- // `buf` isn 't used concurrently.
129+ // `buf` and `didCleanup` aren 't used concurrently.
129130 nonisolated ( unsafe) let buf = UnsafeMutableBufferPointer< UInt8> . allocate( capacity: Int ( getpagesize ( ) ) )
131+ nonisolated ( unsafe) var didCleanup = false
132+
133+ let cleanupRelay : @Sendable ( ) -> Void = {
134+ if didCleanup { return }
135+ didCleanup = true
136+ self . cleanupRelay ( readFd: readFromFd, writeFd: writeToFd, buffer: buf, log: self . log)
137+ }
130138
131139 try ProcessSupervisor . default. poller. add ( readFromFd, mask: EPOLLIN) { mask in
132140 if mask. isHangup && !mask. readyToRead {
133- self . cleanupRelay ( readFd : readFromFd , writeFd : writeToFd , buffer : buf , log : self . log )
141+ cleanupRelay ( )
134142 return
135143 }
136144 // Loop so that in the case that someone wrote > buf.count down the pipe
@@ -146,7 +154,7 @@ final class StandardIO: ManagedProcess.IO & Sendable {
146154 let w = writeTo. write ( view)
147155 if w. wrote != r. read {
148156 self . log? . error ( " stopping relay: short write for stdio " )
149- self . cleanupRelay ( readFd : readFromFd , writeFd : writeToFd , buffer : buf , log : self . log )
157+ cleanupRelay ( )
150158 return
151159 }
152160 }
@@ -156,13 +164,13 @@ final class StandardIO: ManagedProcess.IO & Sendable {
156164 self . log? . error ( " failed with errno \( errno) while reading for fd \( readFromFd) " )
157165 fallthrough
158166 case . eof:
159- self . cleanupRelay ( readFd : readFromFd , writeFd : writeToFd , buffer : buf , log : self . log )
167+ cleanupRelay ( )
160168 self . log? . debug ( " closing relay for \( readFromFd) " )
161169 return
162170 case . again:
163171 // We read all we could, exit.
164172 if mask. isHangup {
165- self . cleanupRelay ( readFd : readFromFd , writeFd : writeToFd , buffer : buf , log : self . log )
173+ cleanupRelay ( )
166174 }
167175 return
168176 default :
0 commit comments